From 2522c79cb3218187063757543fa5747e5acdd081 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 12 Mar 2021 13:36:26 -0500 Subject: [PATCH] Catch async MessageTimeoutException in IP test - fix async exec thread names --- .../ip/tcp/TcpOutboundGatewayTests.java | 2 +- .../ip/tcp/TcpReceivingChannelAdapterTests.java | 4 ++-- .../ip/tcp/TcpSendingMessageHandlerTests.java | 2 +- .../CachingClientConnectionFactoryTests.java | 15 +++++++++++---- .../ip/tcp/connection/ConnectionFactoryTests.java | 2 +- .../FailoverClientConnectionFactoryTests.java | 2 +- .../ip/tcp/connection/TcpNioConnectionTests.java | 2 +- .../ip/tcp/serializer/DeserializationTests.java | 4 ++-- ...atagramPacketMulticastSendingHandlerTests.java | 4 ++-- .../ip/udp/DatagramPacketSendingHandlerTests.java | 6 +++--- .../ip/udp/UdpChannelAdapterTests.java | 2 +- 11 files changed, 26 insertions(+), 19 deletions(-) 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 597c5745fa..950cc032b7 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 @@ -99,7 +99,7 @@ public class TcpOutboundGatewayTests { private static final Log logger = LogFactory.getLog(TcpOutboundGatewayTests.class); - private final AsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("TcpOutboundGatewayTests"); + private final AsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("TcpOutboundGatewayTests-"); @Test void testGoodNetSingle() throws Exception { diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpReceivingChannelAdapterTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpReceivingChannelAdapterTests.java index a0edc8f32e..cb8e5326f4 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpReceivingChannelAdapterTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpReceivingChannelAdapterTests.java @@ -100,7 +100,7 @@ public class TcpReceivingChannelAdapterTests extends AbstractTcpChannelAdapterTe final CountDownLatch latch1 = new CountDownLatch(1); final CountDownLatch latch2 = new CountDownLatch(1); final AtomicBoolean done = new AtomicBoolean(); - new SimpleAsyncTaskExecutor("testNetClientMode").execute(() -> { + new SimpleAsyncTaskExecutor("testNetClientMode-").execute(() -> { try { ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 10); serverSocket.set(server); @@ -412,7 +412,7 @@ public class TcpReceivingChannelAdapterTests extends AbstractTcpChannelAdapterTe handler.setConnectionFactory(scf); TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter(); adapter.setConnectionFactory(scf); - Executor te = new SimpleAsyncTaskExecutor("testNioSingleSharedMany"); + Executor te = new SimpleAsyncTaskExecutor("testNioSingleSharedMany-"); scf.setTaskExecutor(te); scf.start(); QueueChannel channel = new QueueChannel(); 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 377b2f177b..3fc9b3e834 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 @@ -84,7 +84,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest private static final Log logger = LogFactory.getLog(TcpSendingMessageHandlerTests.class); - private final AsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("TcpSendingMessageHandlerTests"); + private final AsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("TcpSendingMessageHandlerTests-"); private void readFully(InputStream is, byte[] buff) throws IOException { for (int i = 0; i < buff.length; i++) { 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 fd3ed25a0a..c47398dbec 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 @@ -65,6 +65,7 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.ApplicationEventPublisher; import org.springframework.core.log.LogAccessor; import org.springframework.core.task.SimpleAsyncTaskExecutor; +import org.springframework.integration.MessageTimeoutException; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.ip.IpHeaders; import org.springframework.integration.ip.tcp.TcpOutboundGateway; @@ -751,13 +752,19 @@ public class CachingClientConnectionFactoryTests { invocation.callRealMethod(); String log = ((Supplier) invocation.getArgument(0)).get(); if (log.startsWith("Response")) { - new SimpleAsyncTaskExecutor("testGatewayRelease") - .execute(() -> gate.handleMessage(new GenericMessage<>("bar"))); + new SimpleAsyncTaskExecutor("testGatewayRelease-") + .execute(() -> { + try { + gate.handleMessage(new GenericMessage<>("bar")); + } + catch (MessageTimeoutException e) { + } + }); // hold up the first thread until the second has added its pending reply - latch.await(20, TimeUnit.SECONDS); + this.latch.await(20, TimeUnit.SECONDS); } else if (log.startsWith("Added")) { - latch.countDown(); + this.latch.countDown(); } return null; } 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 98b71a9025..f9626ff655 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 @@ -228,7 +228,7 @@ public class ConnectionFactoryTests { factory.start(); assertThat(latch1.await(10, TimeUnit.SECONDS)).as("missing info log").isTrue(); // stop on a different thread because it waits for the executor - new SimpleAsyncTaskExecutor("testEarlyClose") + new SimpleAsyncTaskExecutor("testEarlyClose-") .execute(factory::stop); int n = 0; DirectFieldAccessor accessor = new DirectFieldAccessor(factory); 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 30962977d3..16ef17ecf2 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 @@ -609,7 +609,7 @@ public class FailoverClientConnectionFactoryTests { private Holder setupAndStartServers(AbstractServerConnectionFactory server1, AbstractServerConnectionFactory server2) { - Executor exec = new SimpleAsyncTaskExecutor("FailoverClientConnectionFactoryTests"); + Executor exec = new SimpleAsyncTaskExecutor("FailoverClientConnectionFactoryTests-"); server1.setTaskExecutor(exec); server2.setTaskExecutor(exec); server1.setBeanName("server1"); 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 ae8f72f7ed..3356a165a0 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 @@ -109,7 +109,7 @@ public class TcpNioConnectionTests { private final ApplicationEventPublisher nullPublisher = mock(ApplicationEventPublisher.class); - private final AsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("TcpNioConnectionTests"); + private final AsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("TcpNioConnectionTests-"); @Test public void testWriteTimeout(TestInfo testInfo) throws Exception { 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 d80510d46b..b6365f4b35 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 @@ -407,7 +407,7 @@ public class DeserializationTests { // eat SocketTimeoutException. Doesn't matter for this test } }; - Executor exec = new SimpleAsyncTaskExecutor("testTimeoutWhileDecoding"); + Executor exec = new SimpleAsyncTaskExecutor("-"); Message message; @@ -459,7 +459,7 @@ public class DeserializationTests { // eat SocketTimeoutException. Doesn't matter for this test } }; - Executor exec = new SimpleAsyncTaskExecutor("testTimeoutWithRawDeserializerEofIsTerminator"); + Executor exec = new SimpleAsyncTaskExecutor("testTimeoutWithRawDeserializerEofIsTerminator-"); Message message; 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 16d8809f5e..17ab13f2fd 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 @@ -90,7 +90,7 @@ public class DatagramPacketMulticastSendingHandlerTests { e.printStackTrace(); } }; - Executor executor = new SimpleAsyncTaskExecutor("verifySendMulticast"); + Executor executor = new SimpleAsyncTaskExecutor("verifySendMulticast-"); executor.execute(catcher); executor.execute(catcher); assertThat(listening.await(10000, TimeUnit.MILLISECONDS)).isTrue(); @@ -191,7 +191,7 @@ public class DatagramPacketMulticastSendingHandlerTests { e.printStackTrace(); } }; - Executor executor = new SimpleAsyncTaskExecutor("verifySendMulticastWithAcks"); + Executor executor = new SimpleAsyncTaskExecutor("verifySendMulticastWithAcks-"); executor.execute(catcher); executor.execute(catcher); assertThat(listening.await(10000, TimeUnit.MILLISECONDS)).isTrue(); 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 a9ef974709..35267d3c00 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2020 the original author or authors. + * Copyright 2002-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -53,7 +53,7 @@ public class DatagramPacketSendingHandlerTests { final CountDownLatch received = new CountDownLatch(1); final AtomicInteger testPort = new AtomicInteger(); final CountDownLatch listening = new CountDownLatch(1); - new SimpleAsyncTaskExecutor() + new SimpleAsyncTaskExecutor("verifySend-") .execute(() -> { try { DatagramSocket socket = new DatagramSocket(); @@ -93,7 +93,7 @@ public class DatagramPacketSendingHandlerTests { final CountDownLatch listening = new CountDownLatch(1); final CountDownLatch ackListening = new CountDownLatch(1); final CountDownLatch ackSent = new CountDownLatch(1); - new SimpleAsyncTaskExecutor() + new SimpleAsyncTaskExecutor("verifySendWithAck-") .execute(() -> { try { DatagramSocket socket = new DatagramSocket(); 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 52d72a5772..b67427f294 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 @@ -190,7 +190,7 @@ 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 - new SimpleAsyncTaskExecutor() + new SimpleAsyncTaskExecutor("testUnicastReceiverWithReply-") .execute(() -> { DatagramPacket answer = new DatagramPacket(new byte[2000], 2000); try {