From 1c66b090fb106edfeebbe32e80535f95d6782f0a Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 31 Aug 2016 17:04:07 -0400 Subject: [PATCH] Lambdas - IP Module --- .../connection/AbstractConnectionFactory.java | 51 +- .../AbstractServerConnectionFactory.java | 9 +- .../tcp/connection/TcpConnectionSupport.java | 10 +- .../udp/UnicastReceivingChannelAdapter.java | 9 +- .../ip/tcp/TcpInboundGatewayTests.java | 71 +- .../ip/tcp/TcpOutboundGatewayTests.java | 415 ++++----- .../ip/tcp/TcpSendingMessageHandlerTests.java | 829 ++++++++---------- .../CachingClientConnectionFactoryTests.java | 156 +--- .../tcp/connection/ConnectionEventTests.java | 16 +- .../ConnectionFactoryShutDownTests.java | 28 +- .../connection/ConnectionFactoryTests.java | 53 +- .../connection/ConnectionTimeoutTests.java | 61 +- .../FailoverClientConnectionFactoryTests.java | 123 +-- .../ip/tcp/connection/SocketSupportTests.java | 108 +-- .../tcp/connection/TcpNetConnectionTests.java | 21 +- .../connection/TcpNioConnectionReadTests.java | 126 +-- .../tcp/connection/TcpNioConnectionTests.java | 230 ++--- .../TcpNioConnectionWriteTests.java | 138 ++- .../tcp/serializer/DeserializationTests.java | 17 +- ...va => LengthHeaderSerializationTests.java} | 2 +- .../ip/tcp/serializer/SerializationTests.java | 125 ++- ...ramPacketMulticastSendingHandlerTests.java | 138 ++- .../DatagramPacketSendingHandlerTests.java | 74 +- .../integration/ip/udp/MultiClientTests.java | 78 +- .../ip/udp/UdpChannelAdapterTests.java | 23 +- .../integration/ip/util/SocketTestUtils.java | 412 ++++----- 26 files changed, 1355 insertions(+), 1968 deletions(-) rename spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/{LenghtHeaderSerializationTests.java => LengthHeaderSerializationTests.java} (98%) diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractConnectionFactory.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractConnectionFactory.java index 93a48717b2..49c3d4939d 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractConnectionFactory.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractConnectionFactory.java @@ -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")); } } }); diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractServerConnectionFactory.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractServerConnectionFactory.java index cece7f1a6b..b82b044d0a 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractServerConnectionFactory.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractServerConnectionFactory.java @@ -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); diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java index b17a8ffbb8..cb9935e1d6 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java @@ -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); } /** diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/UnicastReceivingChannelAdapter.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/UnicastReceivingChannelAdapter.java index e67fdc1c94..3af9ba2bf8 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/UnicastReceivingChannelAdapter.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/UnicastReceivingChannelAdapter.java @@ -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()) { diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpInboundGatewayTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpInboundGatewayTests.java index 963e0839a3..89dc91cddb 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpInboundGatewayTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpInboundGatewayTests.java @@ -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() { - @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() { - @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(errorMessage)); - } + errorChannel.subscribe(message -> { + MessageChannel replyChannel = (MessageChannel) message.getHeaders().getReplyChannel(); + replyChannel.send(new GenericMessage(errorMessage)); }); gateway.setErrorChannel(errorChannel); scf.start(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpOutboundGatewayTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpOutboundGatewayTests.java index 3a91365197..e1c7c051ca 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpOutboundGatewayTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpOutboundGatewayTests.java @@ -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 = new AtomicReference(); - Executors.newSingleThreadExecutor().execute(new Runnable() { - - @Override - public void run() { - try { - ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 100); - serverSocket.set(server); - latch.countDown(); - List sockets = new ArrayList(); - 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 sockets = new ArrayList(); + 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 = new AtomicReference(); - 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 = new AtomicReference(); - 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[] results = (Future[]) new Future[2]; for (int i = 0; i < 2; i++) { final int j = i; - results[j] = (Executors.newSingleThreadExecutor().submit(new Callable() { - - @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 replies = new HashSet(); @@ -355,44 +334,39 @@ public class TcpOutboundGatewayTests extends LogAdjustingTestSupport { final AtomicReference lastReceived = new AtomicReference(); 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() { - - @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 lastReceived = new AtomicReference(); final CountDownLatch serverLatch = new CountDownLatch(1); - Executors.newSingleThreadExecutor().execute(new Runnable() { - - @Override - public void run() { - List sockets = new ArrayList(); - 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 sockets = new ArrayList(); + 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 sockets = new ArrayList(); - 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 sockets = new ArrayList(); + 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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java index d48d278cbf..0d848949cd 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java @@ -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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 serverSockets = new ArrayList(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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 = new AtomicReference(); 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() { - - @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 { diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java index 96500512e7..8db5e151de 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java @@ -463,21 +463,16 @@ public class CachingClientConnectionFactoryTests { public void gatewayIntegrationTest() throws Exception { final List connectionIds = new ArrayList(); 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("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() { - - @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 factories = new ArrayList(); 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 factories = new ArrayList(); 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 connectionIds = new ArrayList(); - 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> received = new AtomicReference>(); - 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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionEventTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionEventTests.java index 78e012f633..49b537c953 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionEventTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionEventTests.java @@ -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 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); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryShutDownTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryShutDownTests.java index d4873c76e4..d25174ed48 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryShutDownTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryShutDownTests.java @@ -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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryTests.java index 185fe6b941..61dd5c2fa9 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionFactoryTests.java @@ -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() { - @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() { - - @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() { - - @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) { diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionTimeoutTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionTimeoutTests.java index 69e684c6ff..1166feb1dc 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionTimeoutTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/ConnectionTimeoutTests.java @@ -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> reply = new AtomicReference>(); 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 serverConnection = new AtomicReference(); - 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 diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactoryTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactoryTests.java index 4adfc89e1d..a7bce7bf53 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactoryTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactoryTests.java @@ -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() { - @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 message = new GenericMessage("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() { - @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() { - @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() { - @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 connectionId = new AtomicReference(); - 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)); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/SocketSupportTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/SocketSupportTests.java index 019b841737..6d153c4cb6 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/SocketSupportTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/SocketSupportTests.java @@ -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() { - - @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> messages = new ArrayList>(); 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> messages = new ArrayList>(); 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> messages = new ArrayList>(); 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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNetConnectionTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNetConnectionTests.java index fc0c798338..959795cfe4 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNetConnectionTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNetConnectionTests.java @@ -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 log = new AtomicReference(); Log logger = mock(Log.class); - doAnswer(new Answer() { - 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> inboundMessage = new AtomicReference>(); - 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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java index 3a6b9b4e4c..debdb500db 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java @@ -71,15 +71,10 @@ public class TcpNioConnectionReadTests { ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - semaphore.release(); - return false; - } - + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + semaphore.release(); + return false; }); // Fire up the sender. @@ -105,21 +100,16 @@ public class TcpNioConnectionReadTests { ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - 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> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - semaphore.release(); - return false; - } - + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + semaphore.release(); + return false; }); // Fire up the sender. @@ -174,15 +159,10 @@ public class TcpNioConnectionReadTests { ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - semaphore.release(); - return false; - } - + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + semaphore.release(); + return false; }); // Fire up the sender. @@ -206,14 +186,9 @@ public class TcpNioConnectionReadTests { final Semaphore semaphore = new Semaphore(0); final List added = new ArrayList(); final List removed = new ArrayList(); - 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 added = new ArrayList(); final List removed = new ArrayList(); - 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 added = new ArrayList(); final List removed = new ArrayList(); - 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 added = new ArrayList(); final List removed = new ArrayList(); - 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 added = new ArrayList(); final List removed = new ArrayList(); - 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 added = new ArrayList(); final List removed = new ArrayList(); - 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) { diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionTests.java index 59b57c6b8a..dc5b11eba6 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionTests.java @@ -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 = new AtomicReference(); - 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 = new AtomicReference(); - 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 = new AtomicReference(); - 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 fields = new ArrayList(); - 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 future = exec.submit(new Callable() { - @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() { - @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 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 future = exec.submit(new Callable() { - @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() { - @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 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)); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionWriteTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionWriteTests.java index fac442025d..2cf6512ac1 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionWriteTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionWriteTests.java @@ -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); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/DeserializationTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/DeserializationTests.java index c7e4c69a78..180097fd97 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/DeserializationTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/DeserializationTests.java @@ -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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/LenghtHeaderSerializationTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/LengthHeaderSerializationTests.java similarity index 98% rename from spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/LenghtHeaderSerializationTests.java rename to spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/LengthHeaderSerializationTests.java index 1835be682f..83aec1440a 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/LenghtHeaderSerializationTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/LengthHeaderSerializationTests.java @@ -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; diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/SerializationTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/SerializationTests.java index 615672db38..dab7521dfb 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/SerializationTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/serializer/SerializationTests.java @@ -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); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketMulticastSendingHandlerTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketMulticastSendingHandlerTests.java index 16f6f83b05..78faa9c4a8 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketMulticastSendingHandlerTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketMulticastSendingHandlerTests.java @@ -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 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 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); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketSendingHandlerTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketSendingHandlerTests.java index bdae6a8640..1179e4a452 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketSendingHandlerTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/DatagramPacketSendingHandlerTests.java @@ -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 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 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); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/MultiClientTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/MultiClientTests.java index e2a3a5cae7..1b58192f70 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/MultiClientTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/MultiClientTests.java @@ -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(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/UdpChannelAdapterTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/UdpChannelAdapterTests.java index 65567bade7..15f9bfe291 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/UdpChannelAdapterTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/UdpChannelAdapterTests.java @@ -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 receivedMessage = (Message) channel.receive(2000); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java index 76dee43c50..5040aa8ae4 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java @@ -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();