From 790d62bcb1d6d42187df703380c70454cf8152f7 Mon Sep 17 00:00:00 2001 From: Rossen Stoyanchev Date: Tue, 29 Apr 2014 23:01:40 -0400 Subject: [PATCH] Simplify and improve STOMP broker relay int tests --- ...erRelayMessageHandlerIntegrationTests.java | 219 +++++++----------- 1 file changed, 81 insertions(+), 138 deletions(-) diff --git a/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerIntegrationTests.java b/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerIntegrationTests.java index c13fa52a73..4c69baa56a 100644 --- a/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerIntegrationTests.java +++ b/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerIntegrationTests.java @@ -20,8 +20,9 @@ import java.nio.charset.Charset; import java.util.ArrayList; import java.util.Arrays; import java.util.List; -import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.BlockingQueue; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import org.apache.activemq.broker.BrokerService; @@ -38,6 +39,7 @@ import org.springframework.messaging.MessageDeliveryException; import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessagingException; import org.springframework.messaging.StubMessageChannel; +import org.springframework.messaging.simp.SimpMessageHeaderAccessor; import org.springframework.messaging.simp.broker.BrokerAvailabilityEvent; import org.springframework.messaging.simp.SimpMessageType; import org.springframework.messaging.support.ExecutorSubscribableChannel; @@ -64,9 +66,9 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { private ExecutorSubscribableChannel responseChannel; - private ExpectationMatchingMessageHandler responseHandler; + private TestMessageHandler responseHandler; - private ExpectationMatchingEventPublisher eventPublisher; + private TestEventPublisher eventPublisher; private int port; @@ -77,9 +79,9 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { this.port = SocketUtils.findAvailableTcpPort(61613); this.responseChannel = new ExecutorSubscribableChannel(); - this.responseHandler = new ExpectationMatchingMessageHandler(); + this.responseHandler = new TestMessageHandler(); this.responseChannel.subscribe(this.responseHandler); - this.eventPublisher = new ExpectationMatchingEventPublisher(); + this.eventPublisher = new TestEventPublisher(); startActiveMqBroker(); createAndStartRelay(); @@ -103,9 +105,8 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { this.relay.setSystemHeartbeatReceiveInterval(0); this.relay.setSystemHeartbeatSendInterval(0); - this.eventPublisher.expectAvailabilityStatusChanges(true); this.relay.start(); - this.eventPublisher.awaitAndAssert(); + this.eventPublisher.expectBrokerAvailabilityEvent(true); } @After @@ -138,35 +139,37 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { @Test public void publishSubscribe() throws Exception { + logger.debug("Starting test publishSubscribe()"); + String sess1 = "sess1"; String sess2 = "sess2"; MessageExchange conn1 = MessageExchangeBuilder.connect(sess1).build(); MessageExchange conn2 = MessageExchangeBuilder.connect(sess2).build(); - this.responseHandler.expect(conn1, conn2); this.relay.handleMessage(conn1.message); this.relay.handleMessage(conn2.message); - this.responseHandler.awaitAndAssert(); + this.responseHandler.expectMessages(conn1, conn2); String subs1 = "subs1"; String destination = "/topic/test"; MessageExchange subscribe = MessageExchangeBuilder.subscribeWithReceipt(sess1, subs1, destination, "r1").build(); - this.responseHandler.expect(subscribe); - this.relay.handleMessage(subscribe.message); - this.responseHandler.awaitAndAssert(); + this.responseHandler.expectMessages(subscribe); MessageExchange send = MessageExchangeBuilder.send(destination, "foo").andExpectMessage(sess1, subs1).build(); - this.responseHandler.expect(send); - this.relay.handleMessage(send.message); - this.responseHandler.awaitAndAssert(); + this.responseHandler.expectMessages(send); } @Test(expected=MessageDeliveryException.class) public void messageDeliverExceptionIfSystemSessionForwardFails() throws Exception { + + logger.debug("Starting test messageDeliveryExceptionIfSystemSessionForwardFails()"); + stopActiveMqBrokerAndAwait(); + this.eventPublisher.expectBrokerAvailabilityEvent(false); + StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.SEND); this.relay.handleMessage(MessageBuilder.withPayload("test".getBytes()).setHeaders(headers).build()); } @@ -174,58 +177,52 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { @Test public void brokerBecomingUnvailableTriggersErrorFrame() throws Exception { + logger.debug("Starting test brokerBecomingUnvailableTriggersErrorFrame()"); + String sess1 = "sess1"; MessageExchange connect = MessageExchangeBuilder.connect(sess1).build(); - this.responseHandler.expect(connect); - this.relay.handleMessage(connect.message); - this.responseHandler.awaitAndAssert(); - - this.responseHandler.expect(MessageExchangeBuilder.error(sess1).build()); + this.responseHandler.expectMessages(connect); + MessageExchange error = MessageExchangeBuilder.error(sess1).build(); stopActiveMqBrokerAndAwait(); - - this.responseHandler.awaitAndAssert(); + this.eventPublisher.expectBrokerAvailabilityEvent(false); + this.responseHandler.expectMessages(error); } @Test public void brokerAvailabilityEventWhenStopped() throws Exception { - this.eventPublisher.expectAvailabilityStatusChanges(false); + + logger.debug("Starting test brokerAvailabilityEventWhenStopped()"); + stopActiveMqBrokerAndAwait(); - this.eventPublisher.awaitAndAssert(); + this.eventPublisher.expectBrokerAvailabilityEvent(false); } @Test public void relayReconnectsIfBrokerComesBackUp() throws Exception { + logger.debug("Starting test relayReconnectsIfBrokerComesBackUp()"); + String sess1 = "sess1"; MessageExchange conn1 = MessageExchangeBuilder.connect(sess1).build(); - this.responseHandler.expect(conn1); - this.relay.handleMessage(conn1.message); - this.responseHandler.awaitAndAssert(); + this.responseHandler.expectMessages(conn1); String subs1 = "subs1"; String destination = "/topic/test"; - MessageExchange subscribe = - MessageExchangeBuilder.subscribeWithReceipt(sess1, subs1, destination, "r1").build(); - this.responseHandler.expect(subscribe); - + MessageExchange subscribe = MessageExchangeBuilder.subscribeWithReceipt(sess1, subs1, destination, "r1").build(); this.relay.handleMessage(subscribe.message); - this.responseHandler.awaitAndAssert(); - - this.responseHandler.expect(MessageExchangeBuilder.error(sess1).build()); + this.responseHandler.expectMessages(subscribe); + MessageExchange error = MessageExchangeBuilder.error(sess1).build(); stopActiveMqBrokerAndAwait(); + this.responseHandler.expectMessages(error); - this.responseHandler.awaitAndAssert(); + this.eventPublisher.expectBrokerAvailabilityEvent(false); - this.eventPublisher.expectAvailabilityStatusChanges(false); - this.eventPublisher.awaitAndAssert(); - - this.eventPublisher.expectAvailabilityStatusChanges(true); startActiveMqBroker(); - this.eventPublisher.awaitAndAssert(); + this.eventPublisher.expectBrokerAvailabilityEvent(true); // TODO The event publisher assertions show that the broker's back up and the system relay session // has reconnected. We need to decide what we want the reconnect behaviour to be for client relay @@ -236,11 +233,11 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { @Test public void disconnectClosesRelaySessionCleanly() throws Exception { - MessageExchange connect = MessageExchangeBuilder.connect("sess1").build(); - this.responseHandler.expect(connect); + logger.debug("Starting test disconnectClosesRelaySessionCleanly()"); + MessageExchange connect = MessageExchangeBuilder.connect("sess1").build(); this.relay.handleMessage(connect.message); - this.responseHandler.awaitAndAssert(); + this.responseHandler.expectMessages(connect); StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.DISCONNECT); headers.setSessionId("sess1"); @@ -249,79 +246,64 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { Thread.sleep(2000); // Check that we have not received an ERROR as a result of the connection closing - this.responseHandler.awaitAndAssert(); + assertTrue("Unexpected messages: " + this.responseHandler.queue, this.responseHandler.queue.isEmpty()); } - /** - * Handles messages by matching them to expectations including a latch to wait for - * the completion of expected messages. - */ - private static class ExpectationMatchingMessageHandler implements MessageHandler { + private static class TestEventPublisher implements ApplicationEventPublisher { - private final Object monitor = new Object(); + private final BlockingQueue eventQueue = new LinkedBlockingQueue<>(); - private final List expected; - - private final List actual = new ArrayList<>(); - - private final List> unexpected = new ArrayList<>(); - - - public ExpectationMatchingMessageHandler(MessageExchange... expected) { - synchronized (this.monitor) { - this.expected = new CopyOnWriteArrayList<>(expected); + @Override + public void publishEvent(ApplicationEvent event) { + logger.debug("Processing ApplicationEvent " + event); + if (event instanceof BrokerAvailabilityEvent) { + this.eventQueue.add((BrokerAvailabilityEvent) event); } } - public void expect(MessageExchange... expected) { - synchronized (this.monitor) { - this.expected.addAll(Arrays.asList(expected)); - } + public void expectBrokerAvailabilityEvent(boolean isBrokerAvailable) throws InterruptedException { + BrokerAvailabilityEvent event = this.eventQueue.poll(10000, TimeUnit.MILLISECONDS); + assertNotNull("Times out waiting for BrokerAvailabilityEvent[" + isBrokerAvailable + "]", event); + assertEquals(isBrokerAvailable, event.isBrokerAvailable()); } + } - public void awaitAndAssert() throws InterruptedException { - long endTime = System.currentTimeMillis() + 10000; - synchronized (this.monitor) { - while (!this.expected.isEmpty() && System.currentTimeMillis() < endTime) { - this.monitor.wait(500); - } - boolean result = this.expected.isEmpty(); - assertTrue(getAsString(), result && this.unexpected.isEmpty()); - } - } + private static class TestMessageHandler implements MessageHandler { + + private final BlockingQueue> queue = new LinkedBlockingQueue<>(); @Override public void handleMessage(Message message) throws MessagingException { - if (StompHeaderAccessor.wrap(message).getMessageType() != SimpMessageType.HEARTBEAT) { - synchronized(this.monitor) { - for (MessageExchange exch : this.expected) { - if (exch.matchMessage(message)) { - if (exch.isDone()) { - this.expected.remove(exch); - this.actual.add(exch); - if (this.expected.isEmpty()) { - this.monitor.notifyAll(); - } - } - return; - } - } - this.unexpected.add(message); - } + if (SimpMessageType.HEARTBEAT == SimpMessageHeaderAccessor.wrap(message).getMessageType()) { + return; + } + this.queue.add(message); + } + + public void expectMessages(MessageExchange... messageExchanges) throws InterruptedException { + + List expectedMessages = + new ArrayList(Arrays.asList(messageExchanges)); + + while (expectedMessages.size() > 0) { + Message message = this.queue.poll(10000, TimeUnit.MILLISECONDS); + assertNotNull("Timed out waiting for messages, expected [" + expectedMessages + "]", message); + + MessageExchange match = findMatch(expectedMessages, message); + assertNotNull("Unexpected message=" + message + ", expected [" + expectedMessages + "]", match); + + expectedMessages.remove(match); } } - public String getAsString() { - StringBuilder sb = new StringBuilder("\n"); - - synchronized (this.monitor) { - sb.append("UNMATCHED EXPECTATIONS:\n").append(this.expected).append("\n"); - sb.append("MATCHED EXPECTATIONS:\n").append(this.actual).append("\n"); - sb.append("UNEXPECTED MESSAGES:\n").append(this.unexpected).append("\n"); + private MessageExchange findMatch(List expectedMessages, Message message) { + for (MessageExchange exchange : expectedMessages) { + if (exchange.matchMessage(message)) { + return exchange; + } } - - return sb.toString(); + return null; } } @@ -564,43 +546,4 @@ public class StompBrokerRelayMessageHandlerIntegrationTests { } - private static class ExpectationMatchingEventPublisher implements ApplicationEventPublisher { - - private final List expected = new ArrayList<>(); - - private final List actual = new ArrayList<>(); - - private final Object monitor = new Object(); - - - public void expectAvailabilityStatusChanges(Boolean... expected) { - synchronized (this.monitor) { - this.expected.addAll(Arrays.asList(expected)); - } - } - - public void awaitAndAssert() throws InterruptedException { - synchronized(this.monitor) { - long endTime = System.currentTimeMillis() + 60000; - while ((this.expected.size() != this.actual.size()) && (System.currentTimeMillis() < endTime)) { - this.monitor.wait(500); - } - assertEquals(this.expected, this.actual); - } - } - - @Override - public void publishEvent(ApplicationEvent event) { - logger.debug("Processing ApplicationEvent " + event); - if (event instanceof BrokerAvailabilityEvent) { - synchronized(this.monitor) { - this.actual.add(((BrokerAvailabilityEvent) event).isBrokerAvailable()); - if (this.actual.size() == this.expected.size()) { - this.monitor.notifyAll(); - } - } - } - } - } - }