From 4286fc6b825e88be66d148e20bfec8088029620d Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 12 Feb 2015 17:49:00 -0500 Subject: [PATCH] INT-3580: Fix More Stomp Test Race Conditions JIRA: https://jira.spring.io/browse/INT-3580 --- .../client/StompIntegrationTests.java | 45 ++++++++++++------- 1 file changed, 29 insertions(+), 16 deletions(-) diff --git a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/client/StompIntegrationTests.java b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/client/StompIntegrationTests.java index 95cdd202ed..b8ca7c9c6d 100644 --- a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/client/StompIntegrationTests.java +++ b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/client/StompIntegrationTests.java @@ -134,6 +134,9 @@ public class StompIntegrationTests { Message message2 = MessageBuilder.withPayload(5).setHeaders(headers).build(); this.webSocketOutputChannel.send(message); + + waitForSubscribe(); + this.webSocketOutputChannel.send(message2); Message receive = webSocketInputChannel.receive(10000); @@ -157,6 +160,9 @@ public class StompIntegrationTests { Message message2 = MessageBuilder.withPayload(10).setHeaders(headers).build(); this.webSocketOutputChannel.send(message); + + waitForSubscribe(); + this.webSocketOutputChannel.send(message2); Message receive = webSocketInputChannel.receive(10000); @@ -209,22 +215,7 @@ public class StompIntegrationTests { this.webSocketOutputChannel.send(message); - SimpleBrokerMessageHandler serverBrokerMessageHandler = - this.serverContext.getBean("simpleBrokerMessageHandler", SimpleBrokerMessageHandler.class); - - SubscriptionRegistry subscriptionRegistry = serverBrokerMessageHandler.getSubscriptionRegistry(); - - @SuppressWarnings("rawtypes") - Map subscriptions = TestUtils.getPropertyValue(subscriptionRegistry, "subscriptionRegistry.sessions", Map.class); - - int n = 0; - - while (subscriptions.isEmpty() && n++ < 100) { - Thread.sleep(100); - subscriptions = TestUtils.getPropertyValue(subscriptionRegistry, "subscriptionRegistry.sessions", Map.class); - } - - assertTrue("The subscription for the 'user/queue/error' hasn't been registered", n < 100); + waitForSubscribe(); this.webSocketOutputChannel.send(message2); @@ -256,6 +247,9 @@ public class StompIntegrationTests { Message message2 = MessageBuilder.withPayload("Bob").setHeaders(headers).build(); this.webSocketOutputChannel.send(message); + + waitForSubscribe(); + this.webSocketOutputChannel.send(message2); Message receive = webSocketInputChannel.receive(10000); @@ -263,6 +257,25 @@ public class StompIntegrationTests { assertEquals("Hello Bob", receive.getPayload()); } + private void waitForSubscribe() throws InterruptedException { + SimpleBrokerMessageHandler serverBrokerMessageHandler = + this.serverContext.getBean("simpleBrokerMessageHandler", SimpleBrokerMessageHandler.class); + + SubscriptionRegistry subscriptionRegistry = serverBrokerMessageHandler.getSubscriptionRegistry(); + + @SuppressWarnings("rawtypes") + Map subscriptions = TestUtils.getPropertyValue(subscriptionRegistry, "subscriptionRegistry.sessions", Map.class); + + int n = 0; + + while (subscriptions.isEmpty() && n++ < 100) { + Thread.sleep(100); + subscriptions = TestUtils.getPropertyValue(subscriptionRegistry, "subscriptionRegistry.sessions", Map.class); + } + + assertTrue("The subscription for the 'user/queue/error' hasn't been registered", n < 100); + } + @Configuration @EnableIntegration