From b7221938aff417787aa0756e05dbf23884052d57 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Sun, 7 Jun 2015 20:17:01 -0400 Subject: [PATCH] INT-3729: Fix STOMP test sporadic failures JIRA: https://jira.spring.io/browse/INT-3729 --- .../converter/PassThruMessageConverter.java | 3 +- .../configuration/EnableIntegrationTests.java | 90 +------------------ .../stomp/outbound/StompMessageHandler.java | 2 +- .../client/StompServerIntegrationTests.java | 52 ++++++++++- ...annelAdapterWebSocketIntegrationTests.java | 1 + 5 files changed, 58 insertions(+), 90 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/converter/PassThruMessageConverter.java b/spring-integration-core/src/main/java/org/springframework/integration/support/converter/PassThruMessageConverter.java index b6997066b7..7ba6d4f7a1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/converter/PassThruMessageConverter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/converter/PassThruMessageConverter.java @@ -23,12 +23,13 @@ import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; /** - * The simple {@link MessageConverter} implementation which contact is to return + * The simple {@link MessageConverter} implementation which contract is to return * {@link Message} as is for both {@code from/to} operations. *

* It is useful in cases of some protocol implementations (e.g. STOMP), * which is based on the "Spring Messaging Foundation" and the further logic * operates only with {@link Message}s, e.g. Spring Integration Adapters. + * * @author Artem Bilan * @since 4.2 */ diff --git a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java index c6a23f810b..74b5896f98 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java @@ -146,25 +146,9 @@ import reactor.spring.context.config.EnableReactor; @ContextConfiguration(loader = AnnotationConfigContextLoader.class, classes = {EnableIntegrationTests.ContextConfiguration.class, EnableIntegrationTests.ContextConfiguration2.class}) @RunWith(SpringJUnit4ClassRunner.class) -@DirtiesContext +@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_EACH_TEST_METHOD) public class EnableIntegrationTests { - private final Log logger = LogFactory.getLog(EnableIntegrationTests.class); - - @Rule - public TestRule watcher = new TestWatcher() { - - protected void starting(Description description) { - logger.debug("STARTING TEST: " + description.getMethodName()); - } - - @Override - protected void finished(Description description) { - logger.debug("FINISHED TEST: " + description.getMethodName()); - } - - }; - @Autowired private ApplicationContext context; @@ -283,20 +267,6 @@ public class EnableIntegrationTests { @Autowired private MessageChannel controlBusChannel; - private static Level existingLogLevel; - - // The temporal hooks to investigate CI failures - @BeforeClass - public static void setup() { - existingLogLevel = LogManager.getLogger("org.springframework.integration").getLevel(); - LogManager.getLogger("org.springframework.integration").setLevel(Level.DEBUG); - } - - @AfterClass - public static void tearDown() { - LogManager.getLogger("org.springframework.integration").setLevel(existingLogLevel); - } - @Test public void testAnnotatedServiceActivator() { assertEquals(10L, TestUtils.getPropertyValue(this.serviceActivatorEndpoint, "maxMessagesPerPoll")); @@ -337,9 +307,6 @@ public class EnableIntegrationTests { assertEquals(10L, TestUtils.getPropertyValue(trigger, "period")); assertFalse(TestUtils.getPropertyValue(trigger, "fixedRate", Boolean.class)); - // Markers to investigate the failures on CI - logger.debug("----SEND Message to 'input' channel----"); - this.input.send(MessageBuilder.withPayload("Foo").build()); Message interceptedMessage = this.wireTapChannel.receive(10000); @@ -360,8 +327,6 @@ public class EnableIntegrationTests { receive = this.publishedChannel.receive(10000); - logger.debug("----RECEIVE Message from 'publishedChannel' channel----" + receive); - assertNotNull(receive); assertEquals("foo", receive.getPayload()); @@ -417,7 +382,6 @@ public class EnableIntegrationTests { } @Test - @DirtiesContext public void testChangePatterns() { try { this.configurer.setComponentNamePatterns(new String[] {"*"}); @@ -729,10 +693,7 @@ public class EnableIntegrationTests { @Override public Message preSend(Message message, MessageChannel channel) { fbInterceptorCounter().incrementAndGet(); - Message message1 = super.preSend(message, channel); - logger.debug("!!!!'ciFactoryBean': the result of 'preSend' on '" + channel + "' '" + - message1 + "'"); - return message1; + return super.preSend(message, channel); } }; } @@ -846,8 +807,6 @@ public class EnableIntegrationTests { @Override public Message preSend(Message message, MessageChannel channel) { this.invoked.incrementAndGet(); - logger.debug("!!!!'TestChannelInterceptor': the result of 'preSend' on '" + channel + "' '" + - message + "'"); return message; } @@ -947,52 +906,9 @@ public class EnableIntegrationTests { }; } - @Bean(name = IntegrationContextUtils.TASK_SCHEDULER_BEAN_NAME) - public TaskScheduler taskScheduler(BeanFactory beanFactory) { - @SuppressWarnings("serial") - ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler() { - - @Override - public void afterPropertiesSet() { - logger.debug("INITIALIZING taskScheduler..."); - super.afterPropertiesSet(); - } - - @Override - public void destroy() { - logger.debug("DESTROYING taskScheduler..."); - super.destroy(); - } - - }; - - taskScheduler.setPoolSize(20); - taskScheduler.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); - taskScheduler.setThreadNamePrefix("my-task-scheduler-"); - MessagePublishingErrorHandler errorHandler = new MessagePublishingErrorHandler(); - errorHandler.setBeanFactory(beanFactory); - taskScheduler.setErrorHandler(errorHandler); - return taskScheduler; - } - @Bean public PollableChannel publishedChannel() { - return new QueueChannel() { - - @Override - protected boolean doSend(Message message, long timeout) { - logger.debug("---- 'publishedChannel' Thread State: " + Thread.currentThread().isInterrupted()); - logger.debug("---- 'publishedChannel' before 'doSend': " + message); - logger.debug("---- 'publishedChannel' state before: " + getQueueSize() + - ", " + getRemainingCapacity()); - boolean b = super.doSend(message, timeout); - logger.debug("---- 'publishedChannel' after 'doSend': " + b); - logger.debug("---- 'publishedChannel' state after: " + getQueueSize() + - ", " + getRemainingCapacity()); - return b; - } - - }; + return new QueueChannel(); } @Bean diff --git a/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/outbound/StompMessageHandler.java b/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/outbound/StompMessageHandler.java index 60ac124238..2d2b495f7e 100644 --- a/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/outbound/StompMessageHandler.java +++ b/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/outbound/StompMessageHandler.java @@ -41,7 +41,7 @@ import org.springframework.messaging.simp.stomp.StompSessionHandlerAdapter; import org.springframework.util.Assert; /** - * The {@link AbstractMessageHandler} implemntation to send messages to STOMP destinations. + * The {@link AbstractMessageHandler} implementation to send messages to STOMP destinations. * * @author Artem Bilan * @since 4.2 diff --git a/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/client/StompServerIntegrationTests.java b/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/client/StompServerIntegrationTests.java index 01e822d250..b32e2ddad2 100644 --- a/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/client/StompServerIntegrationTests.java +++ b/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/client/StompServerIntegrationTests.java @@ -16,15 +16,22 @@ package org.springframework.integration.stomp.client; +import static org.hamcrest.core.IsInstanceOf.instanceOf; import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThat; + +import java.util.concurrent.Executors; import org.apache.activemq.broker.BrokerService; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Test; +import org.springframework.context.ApplicationEvent; +import org.springframework.context.ApplicationListener; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.Lifecycle; import org.springframework.context.annotation.AnnotationConfigApplicationContext; @@ -33,8 +40,11 @@ import org.springframework.context.annotation.Configuration; import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.event.inbound.ApplicationEventListeningMessageProducer; import org.springframework.integration.stomp.Reactor2TcpStompSessionManager; import org.springframework.integration.stomp.StompSessionManager; +import org.springframework.integration.stomp.event.StompIntegrationEvent; +import org.springframework.integration.stomp.event.StompReceiptEvent; import org.springframework.integration.stomp.inbound.StompInboundChannelAdapter; import org.springframework.integration.stomp.outbound.StompMessageHandler; import org.springframework.integration.support.converter.PassThruMessageConverter; @@ -43,7 +53,11 @@ import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.simp.stomp.Reactor2TcpStompClient; +import org.springframework.messaging.simp.stomp.StompCommand; import org.springframework.messaging.support.GenericMessage; +import org.springframework.scheduling.TaskScheduler; +import org.springframework.scheduling.concurrent.ConcurrentTaskScheduler; +import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.util.SocketUtils; /** @@ -69,6 +83,9 @@ public class StompServerIntegrationTests { activeMQBroker.start(); stompClient = new Reactor2TcpStompClient("127.0.0.1", port); stompClient.setMessageConverter(new PassThruMessageConverter()); + ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); + taskScheduler.afterPropertiesSet(); + stompClient.setTaskScheduler(taskScheduler); } @AfterClass @@ -81,12 +98,29 @@ public class StompServerIntegrationTests { ConfigurableApplicationContext context1 = new AnnotationConfigApplicationContext(ContextConfiguration.class); ConfigurableApplicationContext context2 = new AnnotationConfigApplicationContext(ContextConfiguration.class); + PollableChannel stompEvents1 = context1.getBean("stompEvents", PollableChannel.class); + PollableChannel stompEvents2 = context2.getBean("stompEvents", PollableChannel.class); + PollableChannel stompInputChannel1 = context1.getBean("stompInputChannel", PollableChannel.class); PollableChannel stompInputChannel2 = context2.getBean("stompInputChannel", PollableChannel.class); MessageChannel stompOutputChannel1 = context1.getBean("stompOutputChannel", MessageChannel.class); MessageChannel stompOutputChannel2 = context2.getBean("stompOutputChannel", MessageChannel.class); + Message eventMessage = stompEvents1.receive(10000); + assertNotNull(eventMessage); + assertThat(eventMessage.getPayload(), instanceOf(StompReceiptEvent.class)); + StompReceiptEvent stompReceiptEvent = (StompReceiptEvent) eventMessage.getPayload(); + assertEquals(StompCommand.SUBSCRIBE, stompReceiptEvent.getStompCommand()); + assertEquals("/topic/myTopic", stompReceiptEvent.getDestination()); + + eventMessage = stompEvents2.receive(10000); + assertNotNull(eventMessage); + assertThat(eventMessage.getPayload(), instanceOf(StompReceiptEvent.class)); + stompReceiptEvent = (StompReceiptEvent) eventMessage.getPayload(); + assertEquals(StompCommand.SUBSCRIBE, stompReceiptEvent.getStompCommand()); + assertEquals("/topic/myTopic", stompReceiptEvent.getDestination()); + stompOutputChannel1.send(new GenericMessage("Hello, Client#2!".getBytes())); Message receive11 = stompInputChannel1.receive(10000); @@ -138,7 +172,9 @@ public class StompServerIntegrationTests { @Bean public StompSessionManager stompSessionManager() { - return new Reactor2TcpStompSessionManager(stompClient); + Reactor2TcpStompSessionManager stompSessionManager = new Reactor2TcpStompSessionManager(stompClient); + stompSessionManager.setAutoReceipt(true); + return stompSessionManager; } @Bean @@ -162,6 +198,20 @@ public class StompServerIntegrationTests { return handler; } + @Bean + public PollableChannel stompEvents() { + return new QueueChannel(); + } + + @Bean + @SuppressWarnings("unchecked") + public ApplicationListener stompEventListener() { + ApplicationEventListeningMessageProducer producer = new ApplicationEventListeningMessageProducer(); + producer.setEventTypes(StompIntegrationEvent.class); + producer.setOutputChannel(stompEvents()); + return producer; + } + } } diff --git a/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapterWebSocketIntegrationTests.java b/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapterWebSocketIntegrationTests.java index ccd18c6e3a..308ec0faff 100644 --- a/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapterWebSocketIntegrationTests.java +++ b/spring-integration-stomp/src/test/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapterWebSocketIntegrationTests.java @@ -209,6 +209,7 @@ public class StompInboundChannelAdapterWebSocketIntegrationTests { public WebSocketStompClient stompClient(TaskScheduler taskScheduler) { WebSocketStompClient webSocketStompClient = new WebSocketStompClient(webSocketClient()); webSocketStompClient.setMessageConverter(new MappingJackson2MessageConverter()); + webSocketStompClient.setDefaultHeartbeat(new long[]{100, 100}); webSocketStompClient.setTaskScheduler(taskScheduler); return webSocketStompClient; }