Lambdas - IP Module

This commit is contained in:
Gary Russell
2016-08-31 17:04:07 -04:00
parent 868004e9e4
commit 1c66b090fb
26 changed files with 1355 additions and 1968 deletions

View File

@@ -630,36 +630,33 @@ public abstract class AbstractConnectionFactory extends IntegrationObjectSupport
connection = (TcpNioConnection) key.attachment();
connection.setLastRead(System.currentTimeMillis());
try {
this.taskExecutor.execute(new Runnable() {
@Override
public void run() {
boolean delayed = false;
try {
connection.readPacket();
this.taskExecutor.execute(() -> {
boolean delayed = false;
try {
connection.readPacket();
}
catch (RejectedExecutionException e1) {
delayRead(selector, now, key);
delayed = true;
}
catch (Exception e2) {
if (connection.isOpen()) {
logger.error("Exception on read " +
connection.getConnectionId() + " " +
e2.getMessage());
connection.close();
}
catch (RejectedExecutionException e) {
delayRead(selector, now, key);
delayed = true;
else {
logger.debug("Connection closed");
}
catch (Exception e) {
if (connection.isOpen()) {
logger.error("Exception on read " +
connection.getConnectionId() + " " +
e.getMessage());
connection.close();
}
else {
logger.debug("Connection closed");
}
}
if (!delayed) {
if (key.channel().isOpen()) {
key.interestOps(SelectionKey.OP_READ);
selector.wakeup();
}
if (!delayed) {
if (key.channel().isOpen()) {
key.interestOps(SelectionKey.OP_READ);
selector.wakeup();
}
else {
connection.sendExceptionToListener(new EOFException("Connection is closed"));
}
else {
connection.sendExceptionToListener(new EOFException("Connection is closed"));
}
}
});

View File

@@ -209,14 +209,7 @@ public abstract class AbstractServerConnectionFactory extends AbstractConnection
TaskScheduler taskScheduler = this.getTaskScheduler();
if (taskScheduler != null) {
try {
taskScheduler.schedule(new Runnable() {
@Override
public void run() {
eventPublisher.publishEvent(event);
}
}, new Date());
taskScheduler.schedule((Runnable) () -> eventPublisher.publishEvent(event), new Date());
}
catch (TaskRejectedException e) {
eventPublisher.publishEvent(event);

View File

@@ -34,7 +34,6 @@ import org.springframework.core.serializer.Deserializer;
import org.springframework.core.serializer.Serializer;
import org.springframework.integration.ip.IpHeaders;
import org.springframework.integration.ip.tcp.serializer.AbstractByteArraySerializer;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.util.Assert;
@@ -248,14 +247,7 @@ public abstract class TcpConnectionSupport implements TcpConnection {
*/
public void enableManualListenerRegistration() {
this.manualListenerRegistration = true;
this.listener = new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return getListener().onMessage(message);
}
};
this.listener = message -> getListener().onMessage(message);
}
/**

View File

@@ -161,14 +161,7 @@ public class UnicastReceivingChannelAdapter extends AbstractInternetProtocolRece
Executor taskExecutor = getTaskExecutor();
if (taskExecutor != null) {
try {
taskExecutor.execute(new Runnable() {
@Override
public void run() {
doSend(packet);
}
});
taskExecutor.execute(() -> doSend(packet));
}
catch (RejectedExecutionException e) {
if (logger.isDebugEnabled()) {

View File

@@ -50,10 +50,7 @@ import org.springframework.integration.ip.tcp.connection.TcpNioServerConnectionF
import org.springframework.integration.ip.util.TestingUtilities;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.SubscribableChannel;
import org.springframework.messaging.core.DestinationResolver;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
@@ -76,12 +73,7 @@ public class TcpInboundGatewayTests {
final QueueChannel channel = new QueueChannel();
gateway.setRequestChannel(channel);
ServiceActivatingHandler handler = new ServiceActivatingHandler(new Service());
handler.setChannelResolver(new DestinationResolver<MessageChannel>() {
@Override
public MessageChannel resolveDestination(String channelName) {
return channel;
}
});
handler.setChannelResolver(channelName -> channel);
Socket socket1 = SocketFactory.getDefault().createSocket("localhost", port);
socket1.getOutputStream().write("Test1\r\n".getBytes());
Socket socket2 = SocketFactory.getDefault().createSocket("localhost", port);
@@ -127,30 +119,27 @@ public class TcpInboundGatewayTests {
final CountDownLatch latch2 = new CountDownLatch(1);
final CountDownLatch latch3 = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 10);
port.set(server.getLocalPort());
latch1.countDown();
Socket socket = server.accept();
socket.getOutputStream().write("Test1\r\nTest2\r\n".getBytes());
byte[] bytes = new byte[12];
readFully(socket.getInputStream(), bytes);
assertEquals("Echo:Test1\r\n", new String(bytes));
readFully(socket.getInputStream(), bytes);
assertEquals("Echo:Test2\r\n", new String(bytes));
latch2.await();
socket.close();
server.close();
done.set(true);
latch3.countDown();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 10);
port.set(server.getLocalPort());
latch1.countDown();
Socket socket = server.accept();
socket.getOutputStream().write("Test1\r\nTest2\r\n".getBytes());
byte[] bytes = new byte[12];
readFully(socket.getInputStream(), bytes);
assertEquals("Echo:Test1\r\n", new String(bytes));
readFully(socket.getInputStream(), bytes);
assertEquals("Echo:Test2\r\n", new String(bytes));
latch2.await();
socket.close();
server.close();
done.set(true);
latch3.countDown();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
@@ -196,12 +185,7 @@ public class TcpInboundGatewayTests {
gateway.setRequestChannel(channel);
gateway.setBeanFactory(mock(BeanFactory.class));
ServiceActivatingHandler handler = new ServiceActivatingHandler(new Service());
handler.setChannelResolver(new DestinationResolver<MessageChannel>() {
@Override
public MessageChannel resolveDestination(String channelName) {
return channel;
}
});
handler.setChannelResolver(channelName -> channel);
Socket socket1 = SocketFactory.getDefault().createSocket("localhost", port);
socket1.getOutputStream().write("Test1\r\n".getBytes());
Socket socket2 = SocketFactory.getDefault().createSocket("localhost", port);
@@ -251,12 +235,9 @@ public class TcpInboundGatewayTests {
gateway.setConnectionFactory(scf);
SubscribableChannel errorChannel = new DirectChannel();
final String errorMessage = "An error occurred";
errorChannel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
MessageChannel replyChannel = (MessageChannel) message.getHeaders().getReplyChannel();
replyChannel.send(new GenericMessage<String>(errorMessage));
}
errorChannel.subscribe(message -> {
MessageChannel replyChannel = (MessageChannel) message.getHeaders().getReplyChannel();
replyChannel.send(new GenericMessage<String>(errorMessage));
});
gateway.setErrorChannel(errorChannel);
scf.start();

View File

@@ -41,7 +41,6 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Executors;
@@ -93,32 +92,27 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 100);
serverSocket.set(server);
latch.countDown();
List<Socket> sockets = new ArrayList<Socket>();
int i = 0;
while (true) {
Socket socket = server.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (i++));
sockets.add(socket);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 100);
serverSocket.set(server);
latch.countDown();
List<Socket> sockets = new ArrayList<Socket>();
int i = 0;
while (true) {
Socket socket = server.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (i++));
sockets.add(socket);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
AbstractClientConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -162,30 +156,25 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 10);
serverSocket.set(server);
latch.countDown();
int i = 0;
Socket socket = server.accept();
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (i++));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 10);
serverSocket.set(server);
latch.countDown();
int i = 0;
Socket socket = server.accept();
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (i++));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
AbstractClientConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -222,31 +211,26 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
int i = 0;
Socket socket = server.accept();
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
Thread.sleep(1000);
oos.writeObject("Reply" + (i++));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
int i = 0;
Socket socket = server.accept();
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
Thread.sleep(1000);
oos.writeObject("Reply" + (i++));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
AbstractClientConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -266,14 +250,9 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
Future<Integer>[] results = (Future<Integer>[]) new Future<?>[2];
for (int i = 0; i < 2; i++) {
final int j = i;
results[j] = (Executors.newSingleThreadExecutor().submit(new Callable<Integer>() {
@Override
public Integer call() throws Exception {
gateway.handleMessage(MessageBuilder.withPayload("Test" + j).build());
return 0;
}
results[j] = (Executors.newSingleThreadExecutor().submit(() -> {
gateway.handleMessage(MessageBuilder.withPayload("Test" + j).build());
return 0;
}));
}
Set<String> replies = new HashSet<String>();
@@ -355,44 +334,39 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final AtomicReference<String> lastReceived = new AtomicReference<String>();
final CountDownLatch serverLatch = new CountDownLatch(2);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
latch.countDown();
int i = 0;
while (!done.get()) {
Socket socket = server.accept();
i++;
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request);
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
if (i < 2) {
Thread.sleep(2000);
}
oos.writeObject(request.replace("Test", "Reply"));
logger.debug("Replied to " + request);
lastReceived.set(request);
serverLatch.countDown();
}
catch (IOException e) {
logger.debug("error on write " + e.getClass().getSimpleName());
socket.close();
Executors.newSingleThreadExecutor().execute(() -> {
try {
latch.countDown();
int i = 0;
while (!done.get()) {
Socket socket = server.accept();
i++;
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request);
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
if (i < 2) {
Thread.sleep(2000);
}
oos.writeObject(request.replace("Test", "Reply"));
logger.debug("Replied to " + request);
lastReceived.set(request);
serverLatch.countDown();
}
catch (IOException e1) {
logger.debug("error on write " + e1.getClass().getSimpleName());
socket.close();
}
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
catch (Exception e2) {
if (!done.get()) {
e2.printStackTrace();
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
final TcpOutboundGateway gateway = new TcpOutboundGateway();
@@ -414,14 +388,9 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
for (int i = 0; i < 2; i++) {
final int j = i;
results[j] = (Executors.newSingleThreadExecutor().submit(new Callable<Integer>() {
@Override
public Integer call() throws Exception {
gateway.handleMessage(MessageBuilder.withPayload("Test" + j).build());
return j;
}
results[j] = (Executors.newSingleThreadExecutor().submit(() -> {
gateway.handleMessage(MessageBuilder.withPayload("Test" + j).build());
return j;
}));
}
@@ -463,40 +432,35 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final AtomicBoolean done = new AtomicBoolean();
final CountDownLatch serverLatch = new CountDownLatch(1);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
while (!done.get()) {
Socket socket = server.accept();
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request);
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("bar");
logger.debug("Replied to " + request);
serverLatch.countDown();
}
catch (IOException e) {
logger.debug("error on write " + e.getClass().getSimpleName());
socket.close();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
while (!done.get()) {
Socket socket = server.accept();
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request);
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("bar");
logger.debug("Replied to " + request);
serverLatch.countDown();
}
catch (IOException e1) {
logger.debug("error on write " + e1.getClass().getSimpleName());
socket.close();
}
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
catch (Exception e2) {
if (!done.get()) {
e2.printStackTrace();
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
@@ -548,40 +512,35 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final AtomicBoolean done = new AtomicBoolean();
final CountDownLatch serverLatch = new CountDownLatch(1);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
while (!done.get()) {
Socket socket = server.accept();
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request);
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("bar");
logger.debug("Replied to " + request);
serverLatch.countDown();
}
catch (IOException e) {
logger.debug("error on write " + e.getClass().getSimpleName());
socket.close();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
while (!done.get()) {
Socket socket = server.accept();
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request);
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("bar");
logger.debug("Replied to " + request);
serverLatch.countDown();
}
catch (IOException e1) {
logger.debug("error on write " + e1.getClass().getSimpleName());
socket.close();
}
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
catch (Exception e2) {
if (!done.get()) {
e2.printStackTrace();
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
@@ -701,45 +660,40 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final AtomicReference<String> lastReceived = new AtomicReference<String>();
final CountDownLatch serverLatch = new CountDownLatch(1);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
List<Socket> sockets = new ArrayList<Socket>();
try {
latch.countDown();
while (!done.get()) {
Socket socket = server.accept();
sockets.add(socket);
while (!socket.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request + " closing socket");
socket.close();
lastReceived.set(request);
serverLatch.countDown();
}
catch (IOException e) {
socket.close();
}
Executors.newSingleThreadExecutor().execute(() -> {
List<Socket> sockets = new ArrayList<Socket>();
try {
latch.countDown();
while (!done.get()) {
Socket socket1 = server.accept();
sockets.add(socket1);
while (!socket1.isClosed()) {
try {
ObjectInputStream ois = new ObjectInputStream(socket1.getInputStream());
String request = (String) ois.readObject();
logger.debug("Read " + request + " closing socket");
socket1.close();
lastReceived.set(request);
serverLatch.countDown();
}
catch (IOException e1) {
socket1.close();
}
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
for (Socket socket : sockets) {
try {
socket.close();
}
catch (IOException e) {
}
}
catch (Exception e2) {
if (!done.get()) {
e2.printStackTrace();
}
}
for (Socket socket2 : sockets) {
try {
socket2.close();
}
catch (IOException e3) {
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
final TcpOutboundGateway gateway = new TcpOutboundGateway();
@@ -829,31 +783,26 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport {
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
List<Socket> sockets = new ArrayList<Socket>();
try {
latch.countDown();
while (!done.get()) {
sockets.add(server.accept());
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
for (Socket socket : sockets) {
try {
socket.close();
}
catch (IOException e) {
}
Executors.newSingleThreadExecutor().execute(() -> {
List<Socket> sockets = new ArrayList<Socket>();
try {
latch.countDown();
while (!done.get()) {
sockets.add(server.accept());
}
}
catch (Exception e1) {
if (!done.get()) {
e1.printStackTrace();
}
}
for (Socket socket : sockets) {
try {
socket.close();
}
catch (IOException e2) {
}
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
final TcpOutboundGateway gateway = new TcpOutboundGateway();

View File

@@ -48,8 +48,6 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.junit.Test;
import org.mockito.Mockito;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.context.support.AbstractApplicationContext;
@@ -98,30 +96,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + (++i) + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + (++i) + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -156,30 +149,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + (++i) + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + (++i) + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -226,30 +214,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + (++i) + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + (++i) + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -287,30 +270,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("\u0002Reply" + (++i) + "\u0003").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("\u0002Reply" + (++i) + "\u0003").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -345,30 +323,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("\u0002Reply" + (++i) + "\u0003").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("\u0002Reply" + (++i) + "\u0003").getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -406,33 +379,28 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[8];
readFully(socket.getInputStream(), b);
if (!"\u0000\u0000\u0000\u0004Test".equals(new String(b))) {
throw new RuntimeException("Bad Data");
}
b = ("\u0000\u0000\u0000\u0006Reply" + (++i)).getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[8];
readFully(socket.getInputStream(), b);
if (!"\u0000\u0000\u0000\u0004Test".equals(new String(b))) {
throw new RuntimeException("Bad Data");
}
b = ("\u0000\u0000\u0000\u0006Reply" + (++i)).getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -467,33 +435,28 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[8];
readFully(socket.getInputStream(), b);
if (!"\u0000\u0000\u0000\u0004Test".equals(new String(b))) {
throw new RuntimeException("Bad Data");
}
b = ("\u0000\u0000\u0000\u0006Reply" + (++i)).getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
byte[] b = new byte[8];
readFully(socket.getInputStream(), b);
if (!"\u0000\u0000\u0000\u0004Test".equals(new String(b))) {
throw new RuntimeException("Bad Data");
}
b = ("\u0000\u0000\u0000\u0006Reply" + (++i)).getBytes();
socket.getOutputStream().write(b);
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -531,30 +494,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (++i));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (++i));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -588,30 +546,25 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (++i));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ois.readObject();
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (++i));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -649,31 +602,26 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final CountDownLatch latch = new CountDownLatch(1);
final Semaphore semaphore = new Semaphore(0);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 0; i < 2; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
semaphore.release();
socket.close();
}
server.close();
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 0; i < 2; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
semaphore.release();
socket.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -701,31 +649,26 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final CountDownLatch latch = new CountDownLatch(1);
final Semaphore semaphore = new Semaphore(0);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 0; i < 2; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[8];
readFully(socket.getInputStream(), b);
semaphore.release();
socket.close();
}
server.close();
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 0; i < 2; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[8];
readFully(socket.getInputStream(), b);
semaphore.release();
socket.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -753,32 +696,27 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final CountDownLatch latch = new CountDownLatch(1);
final Semaphore semaphore = new Semaphore(0);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 1; i < 3; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + i + "\r\n").getBytes();
socket.getOutputStream().write(b);
socket.close();
}
server.close();
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 1; i < 3; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + i + "\r\n").getBytes();
socket.getOutputStream().write(b);
socket.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -818,32 +756,27 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final CountDownLatch latch = new CountDownLatch(1);
final Semaphore semaphore = new Semaphore(0);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 1; i < 3; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + i + "\r\n").getBytes();
socket.getOutputStream().write(b);
socket.close();
}
server.close();
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
for (int i = 1; i < 3; i++) {
Socket socket = server.accept();
semaphore.release();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
b = ("Reply" + i + "\r\n").getBytes();
socket.getOutputStream().write(b);
socket.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -885,50 +818,41 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicBoolean done = new AtomicBoolean();
final List<Socket> serverSockets = new ArrayList<Socket>();
final ExecutorService exec = Executors.newCachedThreadPool();
exec.execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 100);
serverSocket.set(server);
latch.countDown();
for (int i = 0; i < 100; i++) {
final Socket socket = server.accept();
serverSockets.add(socket);
final int j = i;
exec.execute(new Runnable() {
@Override
public void run() {
semaphore.release();
byte[] b = new byte[9];
try {
readFully(socket.getInputStream(), b);
b = ("Reply" + j + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
catch (IOException e) {
e.printStackTrace();
}
finally {
try {
socket.close();
}
catch (IOException e) { }
}
exec.execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 100);
serverSocket.set(server);
latch.countDown();
for (int i = 0; i < 100; i++) {
final Socket socket = server.accept();
serverSockets.add(socket);
final int j = i;
exec.execute(() -> {
semaphore.release();
byte[] b = new byte[9];
try {
readFully(socket.getInputStream(), b);
b = ("Reply" + j + "\r\n").getBytes();
socket.getOutputStream().write(b);
}
catch (IOException e1) {
e1.printStackTrace();
}
finally {
try {
socket.close();
}
});
}
server.close();
catch (IOException e2) { }
}
});
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -977,43 +901,38 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
Object in = null;
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
if (i == 0) {
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 0;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
Object in = null;
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
if (i == 0) {
in = ois.readObject();
oos.writeObject("Reply" + (++i));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
}
in = ois.readObject();
oos.writeObject("Reply" + (++i));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -1053,39 +972,34 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 100;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
Object in;
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
if (i == 100) {
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
Thread.sleep(500);
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
int i = 100;
while (true) {
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
Object in;
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
if (i == 100) {
in = ois.readObject();
oos.writeObject("Reply" + (i++));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
Thread.sleep(500);
}
in = ois.readObject();
oos.writeObject("Reply" + (i++));
}
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -1127,39 +1041,34 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
Object in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
oos.writeObject("Reply");
socket.close();
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
Object in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
oos.writeObject("Reply");
socket.close();
server.close();
}
catch (Exception e) {
if (!done.get()) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
@@ -1189,39 +1098,34 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
final CountDownLatch latch = new CountDownLatch(1);
final AtomicBoolean done = new AtomicBoolean();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
int i = 0;
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
Object in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (++i));
socket.close();
server.close();
}
catch (Exception e) {
if (i == 0) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
int i = 0;
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
Object in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
in = ois.readObject();
logger.debug("read object: " + in);
oos.writeObject("world!");
ois = new ObjectInputStream(socket.getInputStream());
oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject("Reply" + (++i));
socket.close();
server.close();
}
catch (Exception e) {
if (i == 0) {
e.printStackTrace();
}
}
});
assertTrue(latch.await(10, TimeUnit.SECONDS));
AbstractConnectionFactory ccf = new TcpNioClientConnectionFactory("localhost",
@@ -1266,13 +1170,8 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
public void testConnectionException() throws Exception {
TcpSendingMessageHandler handler = new TcpSendingMessageHandler();
AbstractConnectionFactory mockCcf = mock(AbstractClientConnectionFactory.class);
Mockito.doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
throw new SocketException("Failed to connect");
}
Mockito.doAnswer(invocation -> {
throw new SocketException("Failed to connect");
}).when(mockCcf).getConnection();
handler.setConnectionFactory(mockCcf);
try {

View File

@@ -463,21 +463,16 @@ public class CachingClientConnectionFactoryTests {
public void gatewayIntegrationTest() throws Exception {
final List<String> connectionIds = new ArrayList<String>();
final AtomicBoolean okToRun = new AtomicBoolean(true);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
while (okToRun.get()) {
Message<?> m = inbound.receive(1000);
if (m != null) {
connectionIds.add((String) m.getHeaders().get(IpHeaders.CONNECTION_ID));
replies.send(MessageBuilder.withPayload("foo:" + new String((byte[]) m.getPayload()))
.copyHeaders(m.getHeaders())
.build());
}
Executors.newSingleThreadExecutor().execute(() -> {
while (okToRun.get()) {
Message<?> m = inbound.receive(1000);
if (m != null) {
connectionIds.add((String) m.getHeaders().get(IpHeaders.CONNECTION_ID));
replies.send(MessageBuilder.withPayload("foo:" + new String((byte[]) m.getPayload()))
.copyHeaders(m.getHeaders())
.build());
}
}
});
TestingUtilities.waitListening(serverCf, null);
toGateway.send(new GenericMessage<String>("Hello, world!"));
@@ -546,14 +541,7 @@ public class CachingClientConnectionFactoryTests {
when(factory1.isActive()).thenReturn(true);
when(factory2.isActive()).thenReturn(true);
doThrow(new IOException("fail")).when(mockConn1).send(Mockito.any(Message.class));
doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
return null;
}
}).when(mockConn2).send(Mockito.any(Message.class));
doAnswer(invocation -> null).when(mockConn2).send(Mockito.any(Message.class));
FailoverClientConnectionFactory failoverFactory = new FailoverClientConnectionFactory(factories);
failoverFactory.start();
@@ -572,13 +560,9 @@ public class CachingClientConnectionFactoryTests {
TcpNetServerConnectionFactory server1 = new TcpNetServerConnectionFactory(0);
server1.setBeanName("server1");
final CountDownLatch latch1 = new CountDownLatch(3);
server1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch1.countDown();
return false;
}
server1.registerListener(message -> {
latch1.countDown();
return false;
});
server1.start();
TestingUtilities.waitListening(server1, 10000L);
@@ -586,13 +570,9 @@ public class CachingClientConnectionFactoryTests {
TcpNetServerConnectionFactory server2 = new TcpNetServerConnectionFactory(0);
server1.setBeanName("server2");
final CountDownLatch latch2 = new CountDownLatch(2);
server2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch2.countDown();
return false;
}
server2.registerListener(message -> {
latch2.countDown();
return false;
});
server2.start();
TestingUtilities.waitListening(server2, 10000L);
@@ -600,22 +580,10 @@ public class CachingClientConnectionFactoryTests {
// Failover
AbstractClientConnectionFactory factory1 = new TcpNetClientConnectionFactory("localhost", port1);
factory1.setBeanName("client1");
factory1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory1.registerListener(message -> false);
AbstractClientConnectionFactory factory2 = new TcpNetClientConnectionFactory("localhost", port2);
factory2.setBeanName("client2");
factory2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory2.registerListener(message -> false);
List<AbstractClientConnectionFactory> factories = new ArrayList<AbstractClientConnectionFactory>();
factories.add(factory1);
factories.add(factory2);
@@ -659,13 +627,9 @@ public class CachingClientConnectionFactoryTests {
TcpNetServerConnectionFactory server1 = new TcpNetServerConnectionFactory(0);
server1.setBeanName("server1");
final CountDownLatch latch1 = new CountDownLatch(3);
server1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch1.countDown();
return false;
}
server1.registerListener(message -> {
latch1.countDown();
return false;
});
server1.start();
TestingUtilities.waitListening(server1, 10000L);
@@ -673,13 +637,9 @@ public class CachingClientConnectionFactoryTests {
TcpNetServerConnectionFactory server2 = new TcpNetServerConnectionFactory(0);
server1.setBeanName("server2");
final CountDownLatch latch2 = new CountDownLatch(2);
server2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch2.countDown();
return false;
}
server2.registerListener(message -> {
latch2.countDown();
return false;
});
server2.start();
TestingUtilities.waitListening(server2, 10000L);
@@ -687,22 +647,10 @@ public class CachingClientConnectionFactoryTests {
// Failover
AbstractClientConnectionFactory factory1 = new TcpNetClientConnectionFactory("junkjunk", port1);
factory1.setBeanName("client1");
factory1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory1.registerListener(message -> false);
AbstractClientConnectionFactory factory2 = new TcpNetClientConnectionFactory("localhost", port2);
factory2.setBeanName("client2");
factory2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory2.registerListener(message -> false);
List<AbstractClientConnectionFactory> factories = new ArrayList<AbstractClientConnectionFactory>();
factories.add(factory1);
factories.add(factory2);
@@ -737,16 +685,11 @@ public class CachingClientConnectionFactoryTests {
final CountDownLatch latch1 = new CountDownLatch(2);
final CountDownLatch latch2 = new CountDownLatch(102);
final List<String> connectionIds = new ArrayList<String>();
in.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
connectionIds.add((String) message.getHeaders().get(IpHeaders.CONNECTION_ID));
latch1.countDown();
latch2.countDown();
return false;
}
in.registerListener(message -> {
connectionIds.add((String) message.getHeaders().get(IpHeaders.CONNECTION_ID));
latch1.countDown();
latch2.countDown();
return false;
});
in.start();
TestingUtilities.waitListening(in, null);
@@ -783,24 +726,19 @@ public class CachingClientConnectionFactoryTests {
final TcpSendingMessageHandler handler = new TcpSendingMessageHandler();
handler.setConnectionFactory(in);
final AtomicInteger count = new AtomicInteger(2);
in.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
if (!(message instanceof ErrorMessage)) {
if (count.decrementAndGet() < 1) {
try {
Thread.sleep(1000);
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
in.registerListener(message -> {
if (!(message instanceof ErrorMessage)) {
if (count.decrementAndGet() < 1) {
try {
Thread.sleep(1000);
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
handler.handleMessage(message);
}
return false;
handler.handleMessage(message);
}
return false;
});
handler.setBeanFactory(mock(BeanFactory.class));
handler.afterPropertiesSet();
@@ -878,16 +816,12 @@ public class CachingClientConnectionFactoryTests {
factory.setApplicationEventPublisher(mock(ApplicationEventPublisher.class));
final CachingClientConnectionFactory cachingFactory = new CachingClientConnectionFactory(factory, 1);
final AtomicReference<Message<?>> received = new AtomicReference<Message<?>>();
cachingFactory.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
if (!(message instanceof ErrorMessage)) {
received.set(message);
latch.countDown();
}
return false;
cachingFactory.registerListener(message -> {
if (!(message instanceof ErrorMessage)) {
received.set(message);
latch.countDown();
}
return false;
});
cachingFactory.start();

View File

@@ -244,13 +244,7 @@ public class ConnectionEventTests {
});
gw.setConnectionFactory(ccf);
DirectChannel requestChannel = new DirectChannel();
requestChannel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
((MessageChannel) message.getHeaders().getReplyChannel()).send(message);
}
});
requestChannel.subscribe(message -> ((MessageChannel) message.getHeaders().getReplyChannel()).send(message));
gw.start();
Message<String> message = MessageBuilder.withPayload("foo")
.setHeader(IpHeaders.CONNECTION_ID, "bar")
@@ -299,13 +293,7 @@ public class ConnectionEventTests {
});
factory.setBeanName("sf");
factory.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory.registerListener(message -> false);
Log logger = spy(TestUtils.getPropertyValue(factory, "logger", Log.class));
doAnswer(new DoesNothing()).when(logger).error(anyString(), any(Throwable.class));
new DirectFieldAccessor(factory).setPropertyValue("logger", logger);

View File

@@ -48,24 +48,20 @@ public class ConnectionFactoryShutDownTests {
Executor executor = factory.getTaskExecutor();
final CountDownLatch latch1 = new CountDownLatch(1);
final CountDownLatch latch2 = new CountDownLatch(1);
executor.execute(new Runnable() {
@Override
public void run() {
latch1.countDown();
try {
while (true) {
factory.getTaskExecutor();
Thread.sleep(100);
}
executor.execute(() -> {
latch1.countDown();
try {
while (true) {
factory.getTaskExecutor();
Thread.sleep(100);
}
catch (MessagingException e) {
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
latch2.countDown();
}
catch (MessagingException e1) {
}
catch (InterruptedException e2) {
Thread.currentThread().interrupt();
}
latch2.countDown();
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
StopWatch watch = new StopWatch();

View File

@@ -46,8 +46,6 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.junit.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.BeanFactory;
@@ -60,7 +58,6 @@ import org.springframework.integration.ip.event.IpIntegrationEvent;
import org.springframework.integration.ip.tcp.TcpReceivingChannelAdapter;
import org.springframework.integration.test.support.LogAdjustingTestSupport;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
@@ -124,13 +121,10 @@ public class ConnectionFactoryTests extends LogAdjustingTestSupport {
serverFactory.setApplicationEventPublisher(publisher);
serverFactory = spy(serverFactory);
final CountDownLatch serverConnectionInitLatch = new CountDownLatch(1);
doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
Object result = invocation.callRealMethod();
serverConnectionInitLatch.countDown();
return result;
}
doAnswer(invocation -> {
Object result = invocation.callRealMethod();
serverConnectionInitLatch.countDown();
return result;
}).when(serverFactory).wrapConnection(any(TcpConnectionSupport.class));
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
scheduler.setPoolSize(10);
@@ -148,12 +142,7 @@ public class ConnectionFactoryTests extends LogAdjustingTestSupport {
assertThat(((TcpConnectionServerListeningEvent) events.get(0)).getPort(), equalTo(serverFactory.getPort()));
int port = serverFactory.getPort();
TcpNetClientConnectionFactory clientFactory = new TcpNetClientConnectionFactory("localhost", port);
clientFactory.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
clientFactory.registerListener(message -> false);
clientFactory.setBeanName("clientFactory");
clientFactory.setApplicationEventPublisher(publisher);
clientFactory.start();
@@ -226,34 +215,20 @@ public class ConnectionFactoryTests extends LogAdjustingTestSupport {
final CountDownLatch latch3 = new CountDownLatch(1);
when(logger.isInfoEnabled()).thenReturn(true);
when(logger.isDebugEnabled()).thenReturn(true);
doAnswer(new Answer<Void>() {
@Override
public Void answer(InvocationOnMock invocation) throws Throwable {
latch1.countDown();
// wait until the stop nulls the channel
latch2.await(10, TimeUnit.SECONDS);
return null;
}
doAnswer(invocation -> {
latch1.countDown();
// wait until the stop nulls the channel
latch2.await(10, TimeUnit.SECONDS);
return null;
}).when(logger).info(contains("Listening"));
doAnswer(new Answer<Void>() {
@Override
public Void answer(InvocationOnMock invocation) throws Throwable {
latch3.countDown();
return null;
}
doAnswer(invocation -> {
latch3.countDown();
return null;
}).when(logger).debug(contains(message));
factory.start();
assertTrue("missing info log", latch1.await(10, TimeUnit.SECONDS));
// stop on a different thread because it waits for the executor
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
factory.stop();
}
});
Executors.newSingleThreadExecutor().execute(() -> factory.stop());
int n = 0;
DirectFieldAccessor accessor = new DirectFieldAccessor(factory);
while (n++ < 200 && accessor.getPropertyValue(property) != null) {

View File

@@ -73,12 +73,7 @@ public class ConnectionTimeoutTests {
server.start();
TestingUtilities.waitListening(server, null);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", server.getPort());
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
client.registerListener(message -> false);
client.setSoTimeout(1000);
CountDownLatch clientCloseLatch = getCloseLatch(client);
setupClientCallback(client);
@@ -122,15 +117,12 @@ public class ConnectionTimeoutTests {
throws Exception, InterruptedException {
final AtomicReference<Message<?>> reply = new AtomicReference<Message<?>>();
final CountDownLatch replyLatch = new CountDownLatch(1);
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
if (!(message instanceof ErrorMessage)) {
reply.set(message);
replyLatch.countDown();
}
return false;
client.registerListener(message -> {
if (!(message instanceof ErrorMessage)) {
reply.set(message);
replyLatch.countDown();
}
return false;
});
client.setSoTimeout(2000);
CountDownLatch clientClosedLatch = getCloseLatch(client);
@@ -166,14 +158,11 @@ public class ConnectionTimeoutTests {
server.start();
TestingUtilities.waitListening(server, null);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", server.getPort());
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
if (!(message instanceof ErrorMessage)) {
reply.set(message);
}
return false;
client.registerListener(message -> {
if (!(message instanceof ErrorMessage)) {
reply.set(message);
}
return false;
});
client.setSoTimeout(2000);
CountDownLatch clientCloseLatch = getCloseLatch(client);
@@ -206,14 +195,11 @@ public class ConnectionTimeoutTests {
server.start();
TestingUtilities.waitListening(server, null);
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", server.getPort());
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
if (!(message instanceof ErrorMessage)) {
reply.set(message);
}
return false;
client.registerListener(message -> {
if (!(message instanceof ErrorMessage)) {
reply.set(message);
}
return false;
});
client.setSoTimeout(1000);
CountDownLatch clientCloseLatch = getCloseLatch(client);
@@ -234,18 +220,15 @@ public class ConnectionTimeoutTests {
private void setupServerCallbacks(AbstractServerConnectionFactory server, final int serverDelay) {
server.setComponentName("serverFactory");
final AtomicReference<TcpConnection> serverConnection = new AtomicReference<TcpConnection>();
server.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
try {
Thread.sleep(serverDelay);
serverConnection.get().send(message);
}
catch (Exception e) {
e.printStackTrace();
}
return false;
server.registerListener(message -> {
try {
Thread.sleep(serverDelay);
serverConnection.get().send(message);
}
catch (Exception e) {
e.printStackTrace();
}
return false;
});
server.registerSender(new TcpSender() {
@Override

View File

@@ -45,8 +45,6 @@ import org.apache.log4j.Level;
import org.junit.Rule;
import org.junit.Test;
import org.mockito.Mockito;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.context.ApplicationEvent;
@@ -63,8 +61,6 @@ import org.springframework.integration.test.util.TestUtils;
import org.springframework.integration.util.SimplePool;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.SubscribableChannel;
import org.springframework.messaging.support.GenericMessage;
@@ -106,12 +102,7 @@ public class FailoverClientConnectionFactoryTests {
when(factory1.isActive()).thenReturn(true);
when(factory2.isActive()).thenReturn(true);
doThrow(new IOException("fail")).when(conn1).send(Mockito.any(Message.class));
doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
return null;
}
}).when(conn2).send(Mockito.any(Message.class));
doAnswer(invocation -> null).when(conn2).send(Mockito.any(Message.class));
FailoverClientConnectionFactory failoverFactory = new FailoverClientConnectionFactory(factories);
failoverFactory.start();
GenericMessage<String> message = new GenericMessage<String>("foo");
@@ -155,15 +146,12 @@ public class FailoverClientConnectionFactoryTests {
when(factory1.isActive()).thenReturn(true);
when(factory2.isActive()).thenReturn(true);
final AtomicBoolean failedOnce = new AtomicBoolean();
doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
if (!failedOnce.get()) {
failedOnce.set(true);
throw new IOException("fail");
}
return null;
doAnswer(invocation -> {
if (!failedOnce.get()) {
failedOnce.set(true);
throw new IOException("fail");
}
return null;
}).when(conn1).send(Mockito.any(Message.class));
doThrow(new IOException("fail")).when(conn2).send(Mockito.any(Message.class));
FailoverClientConnectionFactory failoverFactory = new FailoverClientConnectionFactory(factories);
@@ -199,12 +187,7 @@ public class FailoverClientConnectionFactoryTests {
factories.add(factory1);
factories.add(factory2);
TcpConnectionSupport conn1 = makeMockConnection();
doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
return null;
}
}).when(conn1).send(Mockito.any(Message.class));
doAnswer(invocation -> null).when(conn1).send(Mockito.any(Message.class));
when(factory1.getConnection()).thenThrow(new IOException("fail")).thenReturn(conn1);
when(factory2.getConnection()).thenThrow(new IOException("fail"));
when(factory1.isActive()).thenReturn(true);
@@ -230,14 +213,11 @@ public class FailoverClientConnectionFactoryTests {
when(factory1.isActive()).thenReturn(true);
when(factory2.isActive()).thenReturn(true);
final AtomicInteger failCount = new AtomicInteger();
doAnswer(new Answer<Object>() {
@Override
public Object answer(InvocationOnMock invocation) throws Throwable {
if (failCount.incrementAndGet() < 3) {
throw new IOException("fail");
}
return null;
doAnswer(invocation -> {
if (failCount.incrementAndGet() < 3) {
throw new IOException("fail");
}
return null;
}).when(conn1).send(Mockito.any(Message.class));
doThrow(new IOException("fail")).when(conn2).send(Mockito.any(Message.class));
FailoverClientConnectionFactory failoverFactory = new FailoverClientConnectionFactory(factories);
@@ -319,13 +299,9 @@ public class FailoverClientConnectionFactoryTests {
TcpNetServerConnectionFactory server1 = new TcpNetServerConnectionFactory(0);
server1.setBeanName("server1");
final CountDownLatch latch1 = new CountDownLatch(3);
server1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch1.countDown();
return false;
}
server1.registerListener(message -> {
latch1.countDown();
return false;
});
server1.start();
TestingUtilities.waitListening(server1, 10000L);
@@ -333,35 +309,19 @@ public class FailoverClientConnectionFactoryTests {
TcpNetServerConnectionFactory server2 = new TcpNetServerConnectionFactory(0);
server2.setBeanName("server2");
final CountDownLatch latch2 = new CountDownLatch(2);
server2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch2.countDown();
return false;
}
server2.registerListener(message -> {
latch2.countDown();
return false;
});
server2.start();
TestingUtilities.waitListening(server2, 10000L);
int port2 = server2.getPort();
AbstractClientConnectionFactory factory1 = new TcpNetClientConnectionFactory("localhost", port1);
factory1.setBeanName("client1");
factory1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory1.registerListener(message -> false);
AbstractClientConnectionFactory factory2 = new TcpNetClientConnectionFactory("localhost", port2);
factory2.setBeanName("client2");
factory2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory2.registerListener(message -> false);
// Cache
CachingClientConnectionFactory cachingFactory1 = new CachingClientConnectionFactory(factory1, 2);
cachingFactory1.setBeanName("cache1");
@@ -470,13 +430,9 @@ public class FailoverClientConnectionFactoryTests {
TcpNetServerConnectionFactory server1 = new TcpNetServerConnectionFactory(0);
server1.setBeanName("server1");
final CountDownLatch latch1 = new CountDownLatch(3);
server1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch1.countDown();
return false;
}
server1.registerListener(message -> {
latch1.countDown();
return false;
});
server1.start();
TestingUtilities.waitListening(server1, 10000L);
@@ -484,13 +440,9 @@ public class FailoverClientConnectionFactoryTests {
TcpNetServerConnectionFactory server2 = new TcpNetServerConnectionFactory(0);
server2.setBeanName("server2");
final CountDownLatch latch2 = new CountDownLatch(2);
server2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
latch2.countDown();
return false;
}
server2.registerListener(message -> {
latch2.countDown();
return false;
});
server2.start();
TestingUtilities.waitListening(server2, 10000L);
@@ -498,22 +450,10 @@ public class FailoverClientConnectionFactoryTests {
AbstractClientConnectionFactory factory1 = new TcpNetClientConnectionFactory("junkjunk", port1);
factory1.setBeanName("client1");
factory1.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory1.registerListener(message -> false);
AbstractClientConnectionFactory factory2 = new TcpNetClientConnectionFactory("localhost", port2);
factory2.setBeanName("client2");
factory2.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
factory2.registerListener(message -> false);
// Cache
CachingClientConnectionFactory cachingFactory1 = new CachingClientConnectionFactory(factory1, 2);
@@ -616,12 +556,9 @@ public class FailoverClientConnectionFactoryTests {
gateway1.setConnectionFactory(server1);
SubscribableChannel channel = new DirectChannel();
final AtomicReference<String> connectionId = new AtomicReference<String>();
channel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
connectionId.set((String) message.getHeaders().get(IpHeaders.CONNECTION_ID));
((MessageChannel) message.getHeaders().getReplyChannel()).send(message);
}
channel.subscribe(message -> {
connectionId.set((String) message.getHeaders().get(IpHeaders.CONNECTION_ID));
((MessageChannel) message.getHeaders().getReplyChannel()).send(message);
});
gateway1.setRequestChannel(channel);
gateway1.setBeanFactory(mock(BeanFactory.class));

View File

@@ -42,8 +42,6 @@ import javax.net.ssl.SSLEngine;
import org.junit.Ignore;
import org.junit.Test;
import org.mockito.Mockito;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer;
import org.springframework.integration.ip.util.TestingUtilities;
@@ -97,15 +95,10 @@ public class SocketSupportTests {
when(factory.createServerSocket(0, 5)).thenReturn(serverSocket);
final CountDownLatch latch1 = new CountDownLatch(1);
final CountDownLatch latch2 = new CountDownLatch(1);
when(serverSocket.accept()).thenReturn(socket).then(new Answer<Socket>() {
@Override
public Socket answer(InvocationOnMock invocation) throws Throwable {
latch1.countDown();
latch2.await(10, TimeUnit.SECONDS);
return null;
}
when(serverSocket.accept()).thenReturn(socket).then(invocation -> {
latch1.countDown();
latch2.await(10, TimeUnit.SECONDS);
return null;
});
TcpSocketSupport socketSupport = mock(TcpSocketSupport.class);
@@ -125,14 +118,7 @@ public class SocketSupportTests {
@Test
public void testNioClientAndServer() throws Exception {
TcpNioServerConnectionFactory serverConnectionFactory = new TcpNioServerConnectionFactory(0);
serverConnectionFactory.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
serverConnectionFactory.registerListener(message -> false);
final AtomicInteger ppSocketCountServer = new AtomicInteger();
final AtomicInteger ppServerSocketCountServer = new AtomicInteger();
final CountDownLatch latch = new CountDownLatch(1);
@@ -292,15 +278,10 @@ Certificate fingerprints:
server.setTcpSocketFactorySupport(tcpSocketFactorySupport);
final List<Message<?>> messages = new ArrayList<Message<?>>();
final CountDownLatch latch = new CountDownLatch(1);
server.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messages.add(message);
latch.countDown();
return false;
}
server.registerListener(message -> {
messages.add(message);
latch.countDown();
return false;
});
server.setMapper(new SSLMapper());
server.start();
@@ -330,15 +311,10 @@ Certificate fingerprints:
server.setTcpSocketFactorySupport(serverTcpSocketFactorySupport);
final List<Message<?>> messages = new ArrayList<Message<?>>();
final CountDownLatch latch = new CountDownLatch(1);
server.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messages.add(message);
latch.countDown();
return false;
}
server.registerListener(message -> {
messages.add(message);
latch.countDown();
return false;
});
server.start();
TestingUtilities.waitListening(server, null);
@@ -371,15 +347,10 @@ Certificate fingerprints:
server.setTcpNioConnectionSupport(tcpNioConnectionSupport);
final List<Message<?>> messages = new ArrayList<Message<?>>();
final CountDownLatch latch = new CountDownLatch(1);
server.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messages.add(message);
latch.countDown();
return false;
}
server.registerListener(message -> {
messages.add(message);
latch.countDown();
return false;
});
server.setMapper(new SSLMapper());
server.start();
@@ -387,14 +358,7 @@ Certificate fingerprints:
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", server.getPort());
client.setTcpNioConnectionSupport(tcpNioConnectionSupport);
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
return false;
}
});
client.registerListener(message -> false);
client.start();
TcpConnection connection = client.getConnection();
@@ -418,21 +382,16 @@ Certificate fingerprints:
final CountDownLatch latch = new CountDownLatch(2);
final Replier replier = new Replier();
server.registerSender(replier);
server.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messages.add(message);
try {
replier.send(message);
}
catch (Exception e) {
e.printStackTrace();
}
latch.countDown();
return false;
server.registerListener(message -> {
messages.add(message);
try {
replier.send(message);
}
catch (Exception e) {
e.printStackTrace();
}
latch.countDown();
return false;
});
ByteArrayCrLfSerializer deserializer = new ByteArrayCrLfSerializer();
deserializer.setMaxMessageSize(120000);
@@ -447,15 +406,10 @@ Certificate fingerprints:
new DefaultTcpNioSSLConnectionSupport(clientSslContextSupport);
clientTcpNioConnectionSupport.afterPropertiesSet();
client.setTcpNioConnectionSupport(clientTcpNioConnectionSupport);
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messages.add(message);
latch.countDown();
return false;
}
client.registerListener(message -> {
messages.add(message);
latch.countDown();
return false;
});
client.setDeserializer(deserializer);
client.start();

View File

@@ -33,8 +33,6 @@ import java.util.concurrent.atomic.AtomicReference;
import org.apache.commons.logging.Log;
import org.junit.Test;
import org.mockito.Mockito;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.context.ApplicationEvent;
@@ -78,11 +76,9 @@ public class TcpNetConnectionTests {
connection.setDeserializer(new ByteArrayStxEtxSerializer());
final AtomicReference<Object> log = new AtomicReference<Object>();
Log logger = mock(Log.class);
doAnswer(new Answer<Object>() {
public Object answer(InvocationOnMock invocation) throws Throwable {
log.set(invocation.getArguments()[0]);
return null;
}
doAnswer(invocation -> {
log.set(invocation.getArguments()[0]);
return null;
}).when(logger).error(Mockito.anyString());
DirectFieldAccessor accessor = new DirectFieldAccessor(connection);
accessor.setPropertyValue("logger", logger);
@@ -140,14 +136,11 @@ public class TcpNetConnectionTests {
out.close();
final AtomicReference<Message<?>> inboundMessage = new AtomicReference<Message<?>>();
TcpListener listener = new TcpListener() {
public boolean onMessage(Message<?> message) {
if (!(message instanceof ErrorMessage)) {
inboundMessage.set(message);
}
return false;
TcpListener listener = message1 -> {
if (!(message1 instanceof ErrorMessage)) {
inboundMessage.set(message1);
}
return false;
};
inboundConnection.registerListener(listener);
inboundConnection.run();

View File

@@ -71,15 +71,10 @@ public class TcpNioConnectionReadTests {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
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.
@@ -105,21 +100,16 @@ public class TcpNioConnectionReadTests {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
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(1000);
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
semaphore.release();
return false;
});
int howMany = 2;
@@ -142,15 +132,10 @@ public class TcpNioConnectionReadTests {
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
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.
@@ -174,15 +159,10 @@ public class TcpNioConnectionReadTests {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
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.
@@ -206,14 +186,9 @@ public class TcpNioConnectionReadTests {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
semaphore.release();
return false;
}
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
semaphore.release();
return false;
}, new TcpSender() {
@Override
@@ -248,14 +223,9 @@ public class TcpNioConnectionReadTests {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
semaphore.release();
return false;
}
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
semaphore.release();
return false;
}, new TcpSender() {
@Override
@@ -290,14 +260,9 @@ public class TcpNioConnectionReadTests {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
semaphore.release();
return false;
}
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
semaphore.release();
return false;
}, new TcpSender() {
@Override
@@ -337,14 +302,9 @@ public class TcpNioConnectionReadTests {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
semaphore.release();
return false;
}
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
semaphore.release();
return false;
}, new TcpSender() {
@Override
@@ -381,14 +341,9 @@ public class TcpNioConnectionReadTests {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
semaphore.release();
return false;
}
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
semaphore.release();
return false;
}, new TcpSender() {
@Override
@@ -455,12 +410,9 @@ public class TcpNioConnectionReadTests {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
responses.add(message);
return false;
}
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
responses.add(message);
return false;
}, new TcpSender() {
@Override
public void addNewConnection(TcpConnection connection) {

View File

@@ -49,7 +49,6 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Callable;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
@@ -89,8 +88,6 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.ReflectionUtils.FieldCallback;
import org.springframework.util.ReflectionUtils.FieldFilter;
/**
@@ -117,22 +114,18 @@ public class TcpNioConnectionTests {
final CountDownLatch latch = new CountDownLatch(1);
final CountDownLatch done = new CountDownLatch(1);
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
@SuppressWarnings("unused")
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
logger.debug(testName.getMethodName() + " starting server for " + server.getLocalPort());
serverSocket.set(server);
latch.countDown();
Socket s = server.accept();
// block so we fill the buffer
done.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
logger.debug(testName.getMethodName() + " starting server for " + server.getLocalPort());
serverSocket.set(server);
latch.countDown();
Socket s = server.accept();
// block so we fill the buffer
done.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
@@ -158,23 +151,20 @@ public class TcpNioConnectionTests {
final CountDownLatch latch = new CountDownLatch(1);
final CountDownLatch done = new CountDownLatch(1);
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
logger.debug(testName.getMethodName() + " starting server for " + server.getLocalPort());
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
// block to cause timeout on read.
done.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
logger.debug(testName.getMethodName() + " starting server for " + server.getLocalPort());
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
// block to cause timeout on read.
done.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
@@ -206,21 +196,18 @@ public class TcpNioConnectionTests {
public void testMemoryLeak() throws Exception {
final CountDownLatch latch = new CountDownLatch(1);
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
logger.debug(testName.getMethodName() + " starting server for " + server.getLocalPort());
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
}
catch (Exception e) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0);
logger.debug(testName.getMethodName() + " starting server for " + server.getLocalPort());
serverSocket.set(server);
latch.countDown();
Socket socket = server.accept();
byte[] b = new byte[6];
readFully(socket.getInputStream(), b);
}
catch (Exception e) {
e.printStackTrace();
}
});
assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
@@ -269,21 +256,10 @@ public class TcpNioConnectionTests {
connections.put(chan2, conn2);
connections.put(chan3, conn3);
final List<Field> fields = new ArrayList<Field>();
ReflectionUtils.doWithFields(SocketChannel.class, new FieldCallback() {
@Override
public void doWith(Field field) throws IllegalArgumentException,
IllegalAccessException {
field.setAccessible(true);
fields.add(field);
}
}, new FieldFilter() {
@Override
public boolean matches(Field field) {
return field.getName().equals("open");
}
});
ReflectionUtils.doWithFields(SocketChannel.class, field -> {
field.setAccessible(true);
fields.add(field);
}, field -> field.getName().equals("open"));
Field field = fields.get(0);
// Can't use Mockito because isOpen() is final
ReflectionUtils.setField(field, chan1, true);
@@ -322,38 +298,32 @@ public class TcpNioConnectionTests {
@Test
public void testInsufficientThreads() throws Exception {
final ExecutorService exec = Executors.newFixedThreadPool(2);
Future<Object> future = exec.submit(new Callable<Object>() {
@Override
public Object call() throws Exception {
SocketChannel channel = mock(SocketChannel.class);
Socket socket = mock(Socket.class);
Mockito.when(channel.socket()).thenReturn(socket);
doAnswer(new Answer<Integer>() {
@Override
public Integer answer(InvocationOnMock invocation) throws Throwable {
ByteBuffer buffer = (ByteBuffer) invocation.getArguments()[0];
buffer.position(1);
return 1;
}
}).when(channel).read(Mockito.any(ByteBuffer.class));
when(socket.getReceiveBufferSize()).thenReturn(1024);
final TcpNioConnection connection = new TcpNioConnection(channel, false, false, nullPublisher, null);
connection.setTaskExecutor(exec);
connection.setPipeTimeout(200);
Method method = TcpNioConnection.class.getDeclaredMethod("doRead");
method.setAccessible(true);
// Nobody reading, should timeout on 6th write.
try {
for (int i = 0; i < 6; i++) {
method.invoke(connection);
}
Future<Object> future = exec.submit(() -> {
SocketChannel channel = mock(SocketChannel.class);
Socket socket = mock(Socket.class);
Mockito.when(channel.socket()).thenReturn(socket);
doAnswer(invocation -> {
ByteBuffer buffer = (ByteBuffer) invocation.getArguments()[0];
buffer.position(1);
return 1;
}).when(channel).read(Mockito.any(ByteBuffer.class));
when(socket.getReceiveBufferSize()).thenReturn(1024);
final TcpNioConnection connection = new TcpNioConnection(channel, false, false, nullPublisher, null);
connection.setTaskExecutor(exec);
connection.setPipeTimeout(200);
Method method = TcpNioConnection.class.getDeclaredMethod("doRead");
method.setAccessible(true);
// Nobody reading, should timeout on 6th write.
try {
for (int i = 0; i < 6; i++) {
method.invoke(connection);
}
catch (Exception e) {
e.printStackTrace();
throw (Exception) e.getCause();
}
return null;
}
catch (Exception e) {
e.printStackTrace();
throw (Exception) e.getCause();
}
return null;
});
try {
Object o = future.get(10, TimeUnit.SECONDS);
@@ -368,46 +338,40 @@ public class TcpNioConnectionTests {
public void testSufficientThreads() throws Exception {
final ExecutorService exec = Executors.newFixedThreadPool(3);
final CountDownLatch messageLatch = new CountDownLatch(1);
Future<Object> future = exec.submit(new Callable<Object>() {
@Override
public Object call() throws Exception {
SocketChannel channel = mock(SocketChannel.class);
Socket socket = mock(Socket.class);
Mockito.when(channel.socket()).thenReturn(socket);
doAnswer(new Answer<Integer>() {
@Override
public Integer answer(InvocationOnMock invocation) throws Throwable {
ByteBuffer buffer = (ByteBuffer) invocation.getArguments()[0];
buffer.position(1025);
buffer.put((byte) '\r');
buffer.put((byte) '\n');
return 1027;
}
}).when(channel).read(Mockito.any(ByteBuffer.class));
final TcpNioConnection connection = new TcpNioConnection(channel, false, false, null, null);
connection.setTaskExecutor(exec);
connection.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messageLatch.countDown();
return false;
}
});
connection.setMapper(new TcpMessageMapper());
connection.setDeserializer(new ByteArrayCrLfSerializer());
Method method = TcpNioConnection.class.getDeclaredMethod("doRead");
method.setAccessible(true);
try {
for (int i = 0; i < 20; i++) {
method.invoke(connection);
}
Future<Object> future = exec.submit(() -> {
SocketChannel channel = mock(SocketChannel.class);
Socket socket = mock(Socket.class);
Mockito.when(channel.socket()).thenReturn(socket);
doAnswer(invocation -> {
ByteBuffer buffer = (ByteBuffer) invocation.getArguments()[0];
buffer.position(1025);
buffer.put((byte) '\r');
buffer.put((byte) '\n');
return 1027;
}).when(channel).read(Mockito.any(ByteBuffer.class));
final TcpNioConnection connection = new TcpNioConnection(channel, false, false, null, null);
connection.setTaskExecutor(exec);
connection.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
messageLatch.countDown();
return false;
}
catch (Exception e) {
e.printStackTrace();
throw (Exception) e.getCause();
});
connection.setMapper(new TcpMessageMapper());
connection.setDeserializer(new ByteArrayCrLfSerializer());
Method method = TcpNioConnection.class.getDeclaredMethod("doRead");
method.setAccessible(true);
try {
for (int i = 0; i < 20; i++) {
method.invoke(connection);
}
return null;
}
catch (Exception e) {
e.printStackTrace();
throw (Exception) e.getCause();
}
return null;
});
future.get(60, TimeUnit.SECONDS);
assertTrue(messageLatch.await(10, TimeUnit.SECONDS));

View File

@@ -60,19 +60,16 @@ public class TcpNioConnectionWriteTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(false, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(false, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -96,19 +93,16 @@ public class TcpNioConnectionWriteTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(false, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(false, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -132,19 +126,16 @@ public class TcpNioConnectionWriteTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(false, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(false, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -168,19 +159,16 @@ public class TcpNioConnectionWriteTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(true, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(true, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -204,19 +192,16 @@ public class TcpNioConnectionWriteTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(true, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(true, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -240,19 +225,16 @@ public class TcpNioConnectionWriteTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(true, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
AbstractConnectionFactory ccf = getClientConnectionFactory(true, port, serializer);
TcpConnection connection = ccf.getConnection();
connection.send(MessageBuilder.withPayload(testString.getBytes()).build());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);

View File

@@ -389,18 +389,13 @@ public class DeserializationTests {
out.setBeanFactory(mock(BeanFactory.class));
out.afterPropertiesSet();
out.start();
Runnable command = new Runnable() {
@Override
public void run() {
try {
out.handleMessage(MessageBuilder.withPayload("\u0004Test").build());
}
catch (Exception e) {
// eat SocketTimeoutException. Doesn't matter for this test
}
Runnable command = () -> {
try {
out.handleMessage(MessageBuilder.withPayload("\u0004Test").build());
}
catch (Exception e) {
// eat SocketTimeoutException. Doesn't matter for this test
}
};
ExecutorService exec = Executors.newSingleThreadExecutor();

View File

@@ -33,7 +33,7 @@ import org.junit.Test;
* @since 2.0.4
*
*/
public class LenghtHeaderSerializationTests {
public class LengthHeaderSerializationTests {
private static final String TEST = "Test";
private String test255;

View File

@@ -47,20 +47,17 @@ public class SerializationTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -84,20 +81,17 @@ public class SerializationTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -121,20 +115,17 @@ public class SerializationTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -158,21 +149,18 @@ public class SerializationTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayRawSerializer serializer = new ByteArrayRawSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
socket.close();
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
ByteBuffer buffer = ByteBuffer.allocate(testString.length());
buffer.put(testString.getBytes());
ByteArrayRawSerializer serializer = new ByteArrayRawSerializer();
serializer.serialize(buffer.array(), socket.getOutputStream());
socket.close();
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);
@@ -195,19 +183,16 @@ public class SerializationTests {
final int port = server.getLocalPort();
server.setSoTimeout(10000);
final CountDownLatch latch = new CountDownLatch(1);
Thread t = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
DefaultSerializer serializer = new DefaultSerializer();
serializer.serialize(testString, socket.getOutputStream());
serializer.serialize(testString, socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
Thread t = new Thread(() -> {
try {
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
DefaultSerializer serializer = new DefaultSerializer();
serializer.serialize(testString, socket.getOutputStream());
serializer.serialize(testString, socket.getOutputStream());
latch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
});
t.setDaemon(true);

View File

@@ -64,35 +64,32 @@ public class DatagramPacketMulticastSendingHandlerTests {
final String payload = "foo";
final CountDownLatch listening = new CountDownLatch(2);
final CountDownLatch received = new CountDownLatch(2);
Runnable catcher = new Runnable() {
@Override
public void run() {
try {
byte[] buffer = new byte[8];
DatagramPacket receivedPacket = new DatagramPacket(buffer, buffer.length);
MulticastSocket socket = new MulticastSocket(testPort);
socket.setInterface(InetAddress.getByName(multicastRule.getNic()));
InetAddress group = InetAddress.getByName(multicastAddress);
socket.joinGroup(group);
listening.countDown();
LogFactory.getLog(getClass())
.debug(Thread.currentThread().getName() + " waiting for packet");
socket.receive(receivedPacket);
socket.close();
byte[] src = receivedPacket.getData();
int length = receivedPacket.getLength();
int offset = receivedPacket.getOffset();
byte[] dest = new byte[length];
System.arraycopy(src, offset, dest, 0, length);
assertEquals(payload, new String(dest));
LogFactory.getLog(getClass())
.debug(Thread.currentThread().getName() + " received packet");
received.countDown();
}
catch (Exception e) {
listening.countDown();
e.printStackTrace();
}
Runnable catcher = () -> {
try {
byte[] buffer = new byte[8];
DatagramPacket receivedPacket = new DatagramPacket(buffer, buffer.length);
MulticastSocket socket1 = new MulticastSocket(testPort);
socket1.setInterface(InetAddress.getByName(multicastRule.getNic()));
InetAddress group = InetAddress.getByName(multicastAddress);
socket1.joinGroup(group);
listening.countDown();
LogFactory.getLog(getClass())
.debug(Thread.currentThread().getName() + " waiting for packet");
socket1.receive(receivedPacket);
socket1.close();
byte[] src = receivedPacket.getData();
int length = receivedPacket.getLength();
int offset = receivedPacket.getOffset();
byte[] dest = new byte[length];
System.arraycopy(src, offset, dest, 0, length);
assertEquals(payload, new String(dest));
LogFactory.getLog(getClass())
.debug(Thread.currentThread().getName() + " received packet");
received.countDown();
}
catch (Exception e) {
listening.countDown();
e.printStackTrace();
}
};
Executor executor = Executors.newFixedThreadPool(2);
@@ -127,49 +124,46 @@ public class DatagramPacketMulticastSendingHandlerTests {
final CountDownLatch listening = new CountDownLatch(2);
final CountDownLatch ackListening = new CountDownLatch(1);
final CountDownLatch ackSent = new CountDownLatch(2);
Runnable catcher = new Runnable() {
@Override
public void run() {
try {
byte[] buffer = new byte[1000];
DatagramPacket receivedPacket = new DatagramPacket(buffer, buffer.length);
MulticastSocket socket = new MulticastSocket(testPort);
socket.setInterface(InetAddress.getByName(multicastRule.getNic()));
socket.setSoTimeout(8000);
InetAddress group = InetAddress.getByName(multicastAddress);
socket.joinGroup(group);
listening.countDown();
assertTrue(ackListening.await(10, TimeUnit.SECONDS));
LogFactory.getLog(getClass()).debug(Thread.currentThread().getName() + " waiting for packet");
socket.receive(receivedPacket);
socket.close();
byte[] src = receivedPacket.getData();
int length = receivedPacket.getLength();
int offset = receivedPacket.getOffset();
byte[] dest = new byte[6];
System.arraycopy(src, offset + length - 6, dest, 0, 6);
assertEquals(payload, new String(dest));
LogFactory.getLog(getClass()).debug(Thread.currentThread().getName() + " received packet");
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
mapper.setAcknowledge(true);
mapper.setLengthCheck(true);
Message<byte[]> message = mapper.toMessage(receivedPacket);
Object id = message.getHeaders().get(IpHeaders.ACK_ID);
byte[] ack = id.toString().getBytes();
DatagramPacket ackPack = new DatagramPacket(ack, ack.length,
new InetSocketAddress(multicastRule.getNic(), ackPort.get()));
DatagramSocket out = new DatagramSocket();
out.send(ackPack);
LogFactory.getLog(getClass()).debug(Thread.currentThread().getName() + " sent ack to "
+ ackPack.getSocketAddress());
out.close();
ackSent.countDown();
socket.close();
}
catch (Exception e) {
listening.countDown();
e.printStackTrace();
}
Runnable catcher = () -> {
try {
byte[] buffer = new byte[1000];
DatagramPacket receivedPacket = new DatagramPacket(buffer, buffer.length);
MulticastSocket socket1 = new MulticastSocket(testPort);
socket1.setInterface(InetAddress.getByName(multicastRule.getNic()));
socket1.setSoTimeout(8000);
InetAddress group = InetAddress.getByName(multicastAddress);
socket1.joinGroup(group);
listening.countDown();
assertTrue(ackListening.await(10, TimeUnit.SECONDS));
LogFactory.getLog(getClass()).debug(Thread.currentThread().getName() + " waiting for packet");
socket1.receive(receivedPacket);
socket1.close();
byte[] src = receivedPacket.getData();
int length = receivedPacket.getLength();
int offset = receivedPacket.getOffset();
byte[] dest = new byte[6];
System.arraycopy(src, offset + length - 6, dest, 0, 6);
assertEquals(payload, new String(dest));
LogFactory.getLog(getClass()).debug(Thread.currentThread().getName() + " received packet");
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
mapper.setAcknowledge(true);
mapper.setLengthCheck(true);
Message<byte[]> message = mapper.toMessage(receivedPacket);
Object id = message.getHeaders().get(IpHeaders.ACK_ID);
byte[] ack = id.toString().getBytes();
DatagramPacket ackPack = new DatagramPacket(ack, ack.length,
new InetSocketAddress(multicastRule.getNic(), ackPort.get()));
DatagramSocket out = new DatagramSocket();
out.send(ackPack);
LogFactory.getLog(getClass()).debug(Thread.currentThread().getName() + " sent ack to "
+ ackPack.getSocketAddress());
out.close();
ackSent.countDown();
socket1.close();
}
catch (Exception e) {
listening.countDown();
e.printStackTrace();
}
};
Executor executor = Executors.newFixedThreadPool(2);

View File

@@ -50,20 +50,17 @@ public class DatagramPacketSendingHandlerTests {
final CountDownLatch received = new CountDownLatch(1);
final AtomicInteger testPort = new AtomicInteger();
final CountDownLatch listening = new CountDownLatch(1);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
DatagramSocket socket = new DatagramSocket();
testPort.set(socket.getLocalPort());
listening.countDown();
socket.receive(receivedPacket);
received.countDown();
socket.close();
}
catch (Exception e) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
DatagramSocket socket = new DatagramSocket();
testPort.set(socket.getLocalPort());
listening.countDown();
socket.receive(receivedPacket);
received.countDown();
socket.close();
}
catch (Exception e) {
e.printStackTrace();
}
});
assertTrue(listening.await(10, TimeUnit.SECONDS));
@@ -92,32 +89,29 @@ public class DatagramPacketSendingHandlerTests {
final CountDownLatch listening = new CountDownLatch(1);
final CountDownLatch ackListening = new CountDownLatch(1);
final CountDownLatch ackSent = new CountDownLatch(1);
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
try {
DatagramSocket socket = new DatagramSocket();
testPort.set(socket.getLocalPort());
listening.countDown();
assertTrue(ackListening.await(10, TimeUnit.SECONDS));
socket.receive(receivedPacket);
socket.close();
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
mapper.setAcknowledge(true);
mapper.setLengthCheck(true);
Message<byte[]> message = mapper.toMessage(receivedPacket);
Object id = message.getHeaders().get(IpHeaders.ACK_ID);
byte[] ack = id.toString().getBytes();
DatagramPacket ackPack = new DatagramPacket(ack, ack.length,
new InetSocketAddress("localHost", ackPort.get()));
DatagramSocket out = new DatagramSocket();
out.send(ackPack);
out.close();
ackSent.countDown();
}
catch (Exception e) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
try {
DatagramSocket socket = new DatagramSocket();
testPort.set(socket.getLocalPort());
listening.countDown();
assertTrue(ackListening.await(10, TimeUnit.SECONDS));
socket.receive(receivedPacket);
socket.close();
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
mapper.setAcknowledge(true);
mapper.setLengthCheck(true);
Message<byte[]> message = mapper.toMessage(receivedPacket);
Object id = message.getHeaders().get(IpHeaders.ACK_ID);
byte[] ack = id.toString().getBytes();
DatagramPacket ackPack = new DatagramPacket(ack, ack.length,
new InetSocketAddress("localHost", ackPort.get()));
DatagramSocket out = new DatagramSocket();
out.send(ackPack);
out.close();
ackSent.countDown();
}
catch (Exception e) {
e.printStackTrace();
}
});
listening.await(10000, TimeUnit.MILLISECONDS);

View File

@@ -65,21 +65,18 @@ public class MultiClientTests {
final AtomicBoolean done = new AtomicBoolean();
for (int i = 0; i < drivers; i++) {
Thread t = new Thread(new Runnable() {
@Override
public void run() {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort());
sender.start();
while (true) {
Message<?> message = queueIn.receive();
sender.handleMessage(message);
if (done.get()) {
break;
}
Thread t = new Thread(() -> {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort());
sender.start();
while (true) {
Message<?> message = queueIn.receive();
sender.handleMessage(message);
if (done.get()) {
break;
}
sender.stop();
}
sender.stop();
});
t.setDaemon(true);
t.start();
@@ -114,23 +111,20 @@ public class MultiClientTests {
final AtomicBoolean done = new AtomicBoolean();
for (int i = 0; i < drivers; i++) {
Thread t = new Thread(new Runnable() {
@Override
public void run() {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort(),
false, true, "localhost", 0,
10000);
sender.start();
while (true) {
Message<?> message = queueIn.receive();
sender.handleMessage(message);
if (done.get()) {
break;
}
Thread t = new Thread(() -> {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort(),
false, true, "localhost", 0,
10000);
sender.start();
while (true) {
Message<?> message = queueIn.receive();
sender.handleMessage(message);
if (done.get()) {
break;
}
sender.stop();
}
sender.stop();
});
t.setDaemon(true);
t.start();
@@ -165,24 +159,20 @@ 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(),
true, true, "localhost", 0,
10000);
sender.start();
while (true) {
Message<?> message = queueIn.receive();
sender.handleMessage(message);
if (done.get()) {
break;
}
Thread t = new Thread(() -> {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort(),
true, true, "localhost", 0,
10000);
sender.start();
while (true) {
Message<?> message = queueIn.receive();
sender.handleMessage(message);
if (done.get()) {
break;
}
sender.stop();
}
sender.stop();
});
t.setDaemon(true);
t.start();

View File

@@ -185,19 +185,16 @@ public class UdpChannelAdapterTests {
final CountDownLatch receiverReadyLatch = new CountDownLatch(1);
final CountDownLatch replyReceivedLatch = new CountDownLatch(1);
//main thread sends the reply using the headers, this thread will receive it
Executors.newSingleThreadExecutor().execute(new Runnable() {
@Override
public void run() {
DatagramPacket answer = new DatagramPacket(new byte[2000], 2000);
try {
receiverReadyLatch.countDown();
socket.receive(answer);
theAnswer.set(answer);
replyReceivedLatch.countDown();
}
catch (IOException e) {
e.printStackTrace();
}
Executors.newSingleThreadExecutor().execute(() -> {
DatagramPacket answer = new DatagramPacket(new byte[2000], 2000);
try {
receiverReadyLatch.countDown();
socket.receive(answer);
theAnswer.set(answer);
replyReceivedLatch.countDown();
}
catch (IOException e) {
e.printStackTrace();
}
});
Message<byte[]> receivedMessage = (Message<byte[]>) channel.receive(2000);

View File

@@ -57,38 +57,35 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendLength(final int port, final CountDownLatch latch) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
for (int i = 0; i < 2; i++) {
byte[] len = new byte[4];
ByteBuffer.wrap(len).putInt(TEST_STRING.length() * 2);
socket.getOutputStream().write(len);
socket.getOutputStream().write(TEST_STRING.getBytes());
logger.debug(i + " Wrote first part");
if (latch != null) {
latch.await();
}
Thread.sleep(500);
// send the second chunk
socket.getOutputStream().write(TEST_STRING.getBytes());
logger.debug(i + " Wrote second part");
Thread thread = new Thread(() -> {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
for (int i = 0; i < 2; i++) {
byte[] len = new byte[4];
ByteBuffer.wrap(len).putInt(TEST_STRING.length() * 2);
socket.getOutputStream().write(len);
socket.getOutputStream().write(TEST_STRING.getBytes());
logger.debug(i + " Wrote first part");
if (latch != null) {
latch.await();
}
testCompleteLatch.await(10, TimeUnit.SECONDS);
Thread.sleep(500);
// send the second chunk
socket.getOutputStream().write(TEST_STRING.getBytes());
logger.debug(i + " Wrote second part");
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -102,28 +99,25 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendLengthOverflow(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
byte[] len = new byte[4];
ByteBuffer.wrap(len).putInt(Integer.MAX_VALUE);
socket.getOutputStream().write(len);
socket.getOutputStream().write(TEST_STRING.getBytes());
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
Thread thread = new Thread(() -> {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
byte[] len = new byte[4];
ByteBuffer.wrap(len).putInt(Integer.MAX_VALUE);
socket.getOutputStream().write(len);
socket.getOutputStream().write(TEST_STRING.getBytes());
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -138,34 +132,31 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendFragmented(final int port, final int howMany, final boolean noDelay) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
logger.debug("Connecting to " + port);
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream os = socket.getOutputStream();
for (int i = 0; i < howMany; i++) {
writeByte(os, 0, noDelay);
writeByte(os, 0, noDelay);
writeByte(os, 0, noDelay);
writeByte(os, 2, noDelay);
writeByte(os, 'x', noDelay);
writeByte(os, 'x', noDelay);
}
testCompleteLatch.await(10, TimeUnit.SECONDS);
Thread thread = new Thread(() -> {
Socket socket = null;
try {
logger.debug("Connecting to " + port);
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream os = socket.getOutputStream();
for (int i = 0; i < howMany; i++) {
writeByte(os, 0, noDelay);
writeByte(os, 0, noDelay);
writeByte(os, 0, noDelay);
writeByte(os, 2, noDelay);
writeByte(os, 'x', noDelay);
writeByte(os, 'x', noDelay);
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -189,38 +180,35 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendStxEtx(final int port, final CountDownLatch latch) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 2; i++) {
writeByte(outputStream, 0x02, true);
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote first part");
if (latch != null) {
latch.await();
}
Thread.sleep(500);
// send the second chunk
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote second part");
writeByte(outputStream, 0x03, true);
Thread thread = new Thread(() -> {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 2; i++) {
writeByte(outputStream, 0x02, true);
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote first part");
if (latch != null) {
latch.await();
}
testCompleteLatch.await(10, TimeUnit.SECONDS);
Thread.sleep(500);
// send the second chunk
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote second part");
writeByte(outputStream, 0x03, true);
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -234,29 +222,26 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendStxEtxOverflow(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
writeByte(outputStream, 0x02, true);
for (int i = 0; i < 1500; i++) {
writeByte(outputStream, 'x', true);
}
testCompleteLatch.await(10, TimeUnit.SECONDS);
Thread thread = new Thread(() -> {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
writeByte(outputStream, 0x02, true);
for (int i = 0; i < 1500; i++) {
writeByte(outputStream, 'x', true);
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -271,38 +256,35 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendCrLf(final int port, final CountDownLatch latch) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 2; i++) {
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote first part");
if (latch != null) {
latch.await();
}
Thread.sleep(500);
// send the second chunk
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote second part");
writeByte(outputStream, '\r', true);
writeByte(outputStream, '\n', true);
Thread thread = new Thread(() -> {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 2; i++) {
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote first part");
if (latch != null) {
latch.await();
}
testCompleteLatch.await(10, TimeUnit.SECONDS);
Thread.sleep(500);
// send the second chunk
outputStream.write(TEST_STRING.getBytes());
logger.debug(i + " Wrote second part");
writeByte(outputStream, '\r', true);
writeByte(outputStream, '\n', true);
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -316,24 +298,21 @@ public class SocketTestUtils {
* @param latch Waits for latch to count down before closing the socket.
*/
public static void testSendCrLfSingle(final int port, final CountDownLatch latch) {
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
outputStream.write(TEST_STRING.getBytes());
outputStream.write(TEST_STRING.getBytes());
writeByte(outputStream, '\r', true);
writeByte(outputStream, '\n', true);
if (latch != null) {
latch.await();
}
socket.close();
}
catch (Exception e) {
e.printStackTrace();
Thread thread = new Thread(() -> {
try {
Socket socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
outputStream.write(TEST_STRING.getBytes());
outputStream.write(TEST_STRING.getBytes());
writeByte(outputStream, '\r', true);
writeByte(outputStream, '\n', true);
if (latch != null) {
latch.await();
}
socket.close();
}
catch (Exception e) {
e.printStackTrace();
}
});
thread.setDaemon(true);
@@ -344,19 +323,16 @@ public class SocketTestUtils {
* Sends a single message in two chunks and then closes the socket.
*/
public static void testSendRaw(final int port) {
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
outputStream.write(TEST_STRING.getBytes());
outputStream.write(TEST_STRING.getBytes());
socket.close();
}
catch (Exception e) {
e.printStackTrace();
}
Thread thread = new Thread(() -> {
try {
Socket socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
outputStream.write(TEST_STRING.getBytes());
outputStream.write(TEST_STRING.getBytes());
socket.close();
}
catch (Exception e) {
e.printStackTrace();
}
});
thread.setDaemon(true);
@@ -368,31 +344,28 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendSerialized(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(outputStream);
oos.writeObject(TEST_STRING);
oos.flush();
oos = new ObjectOutputStream(outputStream);
oos.writeObject(TEST_STRING);
oos.flush();
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e) {
e.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e) { }
Thread thread = new Thread(() -> {
Socket socket = null;
try {
socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(outputStream);
oos.writeObject(TEST_STRING);
oos.flush();
oos = new ObjectOutputStream(outputStream);
oos.writeObject(TEST_STRING);
oos.flush();
testCompleteLatch.await(10, TimeUnit.SECONDS);
}
catch (Exception e1) {
e1.printStackTrace();
}
finally {
if (socket != null) {
try {
socket.close();
}
catch (IOException e2) { }
}
}
});
@@ -406,20 +379,17 @@ public class SocketTestUtils {
*/
public static CountDownLatch testSendCrLfOverflow(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(new Runnable() {
@Override
public void run() {
try {
Socket socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 1500; i++) {
writeByte(outputStream, 'x', true);
}
testCompleteLatch.await(10, TimeUnit.SECONDS);
socket.close();
Thread thread = new Thread(() -> {
try {
Socket socket = new Socket(InetAddress.getByName("localhost"), port);
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 1500; i++) {
writeByte(outputStream, 'x', true);
}
catch (Exception e) { }
testCompleteLatch.await(10, TimeUnit.SECONDS);
socket.close();
}
catch (Exception e) { }
});
thread.setDaemon(true);
thread.start();