Remove CONNECT-related message buffer from STOMP relay

Before this change, the StompProtocolHandler always responded to
clients with a CONNECTED frame, while the STOMP broker relay
independantly forwarded the client CONNECT to the broker and waited
for the CONNECTED frame back. That meant the relay had to buffer
client messages until it received the CONNECTED response from
the message broker.

This change ensures that clients wait for a CONNECTED frame from
the message broker. The broker relay forwards the CONNECT frame to
the broker. The broker responds with a CONNECTED frame, which the
relay then forwards to the client. As a result, a (well-written)
client will not send any messages to the relay until the connection
to the broker is fully established.

The StompProtcolHandler can now be configured whether to send CONNECTED
frame back. By default that is off. So when using the simple broker,
the StompProtocolHandler can still respond with CONNECTED frames.

The relay's handling of a connection being dropped has also been
improved. When a connection for a client relay session is dropped
an ERROR frame will be sent back to the client. If a connection is
closed as part of a DISCONNECT frame being sent, no ERROR frame
is sent back to the client. When the connection for the system relay
session is dropped, an event is published indicating that the broker
is unavailable. Reactor's TcpClient will then attempt to re-restablish
the connection.
This commit is contained in:
Andy Wilkinson
2013-09-02 09:45:36 +01:00
committed by Rossen Stoyanchev
parent a489c2cf38
commit 8d2a376b0f
7 changed files with 232 additions and 150 deletions

View File

@@ -53,7 +53,7 @@ public class ServletStompEndpointRegistryTests {
this.webSocketHandler = new SubProtocolWebSocketHandler(channel);
this.queueSuffixResolver = new SimpleUserQueueSuffixResolver();
TaskScheduler taskScheduler = Mockito.mock(TaskScheduler.class);
this.registry = new ServletStompEndpointRegistry(webSocketHandler, queueSuffixResolver, taskScheduler);
this.registry = new ServletStompEndpointRegistry(webSocketHandler, queueSuffixResolver, taskScheduler, false);
}

View File

@@ -30,7 +30,6 @@ import org.apache.commons.logging.LogFactory;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.context.ApplicationEvent;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.messaging.Message;
@@ -63,16 +62,14 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
private ExpectationMatchingEventPublisher eventPublisher;
private int port;
@Before
public void setUp() throws Exception {
int port = SocketUtils.findAvailableTcpPort(61613);
this.port = SocketUtils.findAvailableTcpPort(61613);
this.activeMQBroker = new BrokerService();
this.activeMQBroker.addConnector("stomp://localhost:" + port);
this.activeMQBroker.setStartAsync(false);
this.activeMQBroker.setDeleteAllMessagesOnStartup(true);
this.activeMQBroker.start();
createAndStartBroker();
this.responseChannel = new ExecutorSubscribableChannel();
this.responseHandler = new ExpectationMatchingMessageHandler();
@@ -86,6 +83,14 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
this.relay.start();
}
private void createAndStartBroker() throws Exception {
this.activeMQBroker = new BrokerService();
this.activeMQBroker.addConnector("stomp://localhost:" + port);
this.activeMQBroker.setStartAsync(false);
this.activeMQBroker.setDeleteAllMessagesOnStartup(true);
this.activeMQBroker.start();
}
@After
public void tearDown() throws Exception {
try {
@@ -102,22 +107,24 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
String sess1 = "sess1";
MessageExchange conn1 = MessageExchangeBuilder.connect(sess1).build();
this.relay.handleMessage(conn1.message);
this.responseHandler.expect(conn1);
String sess2 = "sess2";
MessageExchange conn2 = MessageExchangeBuilder.connect(sess2).build();
this.relay.handleMessage(conn2.message);
this.responseHandler.expect(conn2);
this.responseHandler.awaitAndAssert();
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.expect(subscribe);
this.responseHandler.awaitAndAssert();
MessageExchange send = MessageExchangeBuilder.send(destination, "foo").andExpectMessage(sess1, subs1).build();
this.responseHandler.reset();
this.responseHandler.expect(send);
this.relay.handleMessage(send.message);
@@ -129,7 +136,7 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
stopBrokerAndAwait();
MessageExchange connect = MessageExchangeBuilder.connect("sess1").andExpectError().build();
MessageExchange connect = MessageExchangeBuilder.connectWithError("sess1").build();
this.responseHandler.expect(connect);
this.relay.handleMessage(connect.message);
@@ -137,37 +144,31 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
}
@Test
public void brokerUnvailableErrorFrameOnSend() throws Exception {
public void brokerBecomingUnvailableTriggersErrorFrame() throws Exception {
String sess1 = "sess1";
MessageExchange connect = MessageExchangeBuilder.connect(sess1).build();
this.responseHandler.expect(connect);
this.relay.handleMessage(connect.message);
// TODO: expect CONNECTED
Thread.sleep(2000);
this.responseHandler.awaitAndAssert();
this.responseHandler.expect(MessageExchangeBuilder.error(sess1).build());
stopBrokerAndAwait();
MessageExchange subscribe = MessageExchangeBuilder.subscribe(sess1, "s1", "/topic/a").andExpectError().build();
this.responseHandler.expect(subscribe);
this.relay.handleMessage(subscribe.message);
this.responseHandler.awaitAndAssert();
}
@Test
public void brokerAvailabilityEvents() throws Exception {
// TODO: expect CONNECTED
Thread.sleep(2000);
this.eventPublisher.expect(true, false);
this.eventPublisher.expect(true);
this.eventPublisher.awaitAndAssert();
this.eventPublisher.expect(false);
stopBrokerAndAwait();
// TODO: remove when stop is detecteded
this.relay.handleMessage(MessageExchangeBuilder.connect("sess1").build().message);
this.eventPublisher.awaitAndAssert();
}
@@ -176,37 +177,55 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
String sess1 = "sess1";
MessageExchange conn1 = MessageExchangeBuilder.connect(sess1).build();
this.responseHandler.expect(conn1);
this.relay.handleMessage(conn1.message);
this.responseHandler.awaitAndAssert();
String subs1 = "subs1";
String destination = "/topic/test";
MessageExchange subscribe = MessageExchangeBuilder.subscribeWithReceipt(sess1, subs1, destination, "r1").build();
MessageExchange subscribe =
MessageExchangeBuilder.subscribeWithReceipt(sess1, subs1, destination, "r1").build();
this.responseHandler.expect(subscribe);
this.relay.handleMessage(subscribe.message);
this.responseHandler.awaitAndAssert();
this.responseHandler.expect(MessageExchangeBuilder.error(sess1).build());
stopBrokerAndAwait();
// 1st message will see ERROR frame (broker shutdown is not but should be detected)
// 2nd message will be queued (a side effect of CONNECT/CONNECTED-buffering, likely to be removed)
// Finish this once the above changes are made.
this.responseHandler.awaitAndAssert();
this.eventPublisher.expect(true, false);
this.eventPublisher.awaitAndAssert();
this.eventPublisher.expect(true);
createAndStartBroker();
this.eventPublisher.awaitAndAssert();
// 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
// sessions and add further message sending and assertions as appropriate. At the moment any client
// sessions will be closed and an ERROR from will be sent.
}
@Test
public void disconnectClosesRelaySessionCleanly() throws Exception {
String sess1 = "sess1";
MessageExchange conn1 = MessageExchangeBuilder.connect(sess1).build();
this.responseHandler.expect(conn1);
this.relay.handleMessage(conn1.message);
this.responseHandler.awaitAndAssert();
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.DISCONNECT);
headers.setSessionId(sess1);
this.relay.handleMessage(MessageBuilder.withPayloadAndHeaders(new byte[0], headers).build());
/* MessageExchange send = MessageExchangeBuilder.send(destination, "foo").build();
this.responseHandler.reset();
this.relay.handleMessage(send.message);
Thread.sleep(2000);
this.activeMQBroker.start();
Thread.sleep(5000);
send = MessageExchangeBuilder.send(destination, "foo").andExpectMessage(sess1, subs1).build();
this.responseHandler.reset();
this.responseHandler.expect(send);
this.relay.handleMessage(send.message);
// Check that we have not received an ERROR as a result of the connection closing
this.responseHandler.awaitAndAssert();
*/
}
@@ -234,58 +253,66 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
*/
private static class ExpectationMatchingMessageHandler implements MessageHandler {
private final Object monitor = new Object();
private final List<MessageExchange> expected;
private final List<MessageExchange> actual = new CopyOnWriteArrayList<>();
private final List<MessageExchange> actual = new ArrayList<>();
private final List<Message<?>> unexpected = new CopyOnWriteArrayList<>();
private CountDownLatch latch = new CountDownLatch(1);
private final List<Message<?>> unexpected = new ArrayList<>();
public ExpectationMatchingMessageHandler(MessageExchange... expected) {
this.expected = new CopyOnWriteArrayList<>(expected);
synchronized (this.monitor) {
this.expected = new CopyOnWriteArrayList<>(expected);
}
}
public void expect(MessageExchange... expected) {
this.expected.addAll(Arrays.asList(expected));
synchronized (this.monitor) {
this.expected.addAll(Arrays.asList(expected));
}
}
public void awaitAndAssert() throws InterruptedException {
boolean result = this.latch.await(10000, TimeUnit.MILLISECONDS);
assertTrue(getAsString(), result && this.unexpected.isEmpty());
}
public void reset() {
this.latch = new CountDownLatch(1);
this.expected.clear();
this.actual.clear();
this.unexpected.clear();
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());
}
}
@Override
public void handleMessage(Message<?> message) throws MessagingException {
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.latch.countDown();
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;
}
return;
}
this.unexpected.add(message);
}
this.unexpected.add(message);
}
public String getAsString() {
StringBuilder sb = new StringBuilder("\n");
sb.append("INCOMPLETE:\n").append(this.expected).append("\n");
sb.append("COMPLETE:\n").append(this.actual).append("\n");
sb.append("UNMATCHED MESSAGES:\n").append(this.unexpected).append("\n");
synchronized (this.monitor) {
sb.append("INCOMPLETE:\n").append(this.expected).append("\n");
sb.append("COMPLETE:\n").append(this.actual).append("\n");
sb.append("UNMATCHED MESSAGES:\n").append(this.unexpected).append("\n");
}
return sb.toString();
}
}
@@ -352,22 +379,28 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
this.headers = StompHeaderAccessor.wrap(message);
}
public static MessageExchangeBuilder error(String sessionId) {
return new MessageExchangeBuilder(null).andExpectError(sessionId);
}
public static MessageExchangeBuilder connect(String sessionId) {
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.CONNECT);
headers.setSessionId(sessionId);
headers.setAcceptVersion("1.1,1.2");
Message<?> message = MessageBuilder.withPayloadAndHeaders(new byte[0], headers).build();
return new MessageExchangeBuilder(message);
MessageExchangeBuilder builder = new MessageExchangeBuilder(message);
builder.expected.add(new StompConnectedFrameMessageMatcher(sessionId));
return builder;
}
public static MessageExchangeBuilder subscribe(String sessionId, String subscriptionId, String destination) {
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.SUBSCRIBE);
public static MessageExchangeBuilder connectWithError(String sessionId) {
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.CONNECT);
headers.setSessionId(sessionId);
headers.setSubscriptionId(subscriptionId);
headers.setDestination(destination);
headers.setAcceptVersion("1.1,1.2");
Message<?> message = MessageBuilder.withPayloadAndHeaders(new byte[0], headers).build();
return new MessageExchangeBuilder(message);
MessageExchangeBuilder builder = new MessageExchangeBuilder(message);
return builder.andExpectError();
}
public static MessageExchangeBuilder subscribeWithReceipt(String sessionId, String subscriptionId,
@@ -515,35 +548,48 @@ public class StompBrokerRelayMessageHandlerIntegrationTests {
}
}
private static class StompConnectedFrameMessageMatcher extends StompFrameMessageMatcher {
public StompConnectedFrameMessageMatcher(String sessionId) {
super(StompCommand.CONNECTED, sessionId);
}
}
private static class ExpectationMatchingEventPublisher implements ApplicationEventPublisher {
private final List<Boolean> expected = new CopyOnWriteArrayList<>();
private final List<Boolean> expected = new ArrayList<>();
private final List<Boolean> actual = new CopyOnWriteArrayList<>();
private final List<Boolean> actual = new ArrayList<>();
private CountDownLatch latch = new CountDownLatch(1);
private final Object monitor = new Object();
public void expect(Boolean... expected) {
this.expected.addAll(Arrays.asList(expected));
synchronized (this.monitor) {
this.expected.addAll(Arrays.asList(expected));
}
}
public void awaitAndAssert() throws InterruptedException {
if (this.expected.size() == this.actual.size()) {
synchronized(this.monitor) {
long endTime = System.currentTimeMillis() + 5000;
while (this.expected.size() != this.actual.size() && System.currentTimeMillis() < endTime) {
this.monitor.wait(500);
}
assertEquals(this.expected, this.actual);
}
else {
assertTrue("Expected=" + this.expected + ", actual=" + this.actual,
this.latch.await(5, TimeUnit.SECONDS));
}
}
@Override
public void publishEvent(ApplicationEvent event) {
if (event instanceof BrokerAvailabilityEvent) {
this.actual.add(((BrokerAvailabilityEvent) event).isBrokerAvailable());
if (this.actual.size() == this.expected.size()) {
this.latch.countDown();
synchronized(this.monitor) {
this.actual.add(((BrokerAvailabilityEvent) event).isBrokerAvailable());
if (this.actual.size() == this.expected.size()) {
this.monitor.notifyAll();
}
}
}
}

View File

@@ -61,7 +61,33 @@ public class StompProtocolHandlerTests {
}
@Test
public void handleConnect() {
public void connectedResponseIsSentWhenHandlingConnect() {
this.stompHandler.setHandleConnect(true);
TextMessage textMessage = StompTextMessageBuilder.create(StompCommand.CONNECT).headers(
"login:guest", "passcode:guest", "accept-version:1.1,1.0", "heart-beat:10000,10000").build();
this.stompHandler.handleMessageFromClient(this.session, textMessage, this.channel);
verifyNoMoreInteractions(this.channel);
// Check CONNECTED reply
assertEquals(1, this.session.getSentMessages().size());
textMessage = (TextMessage) this.session.getSentMessages().get(0);
Message<?> message = new StompDecoder().decode(ByteBuffer.wrap(textMessage.getPayload().getBytes()));
StompHeaderAccessor replyHeaders = StompHeaderAccessor.wrap(message);
assertEquals(StompCommand.CONNECTED, replyHeaders.getCommand());
assertEquals("1.1", replyHeaders.getVersion());
assertArrayEquals(new long[] {0, 0}, replyHeaders.getHeartbeat());
assertEquals("joe", replyHeaders.getNativeHeader("user-name").get(0));
assertEquals("s1", replyHeaders.getNativeHeader("queue-suffix").get(0));
}
@Test
public void connectIsForwardedWhenNotHandlingConnect() {
this.stompHandler.setHandleConnect(false);
TextMessage textMessage = StompTextMessageBuilder.create(StompCommand.CONNECT).headers(
"login:guest", "passcode:guest", "accept-version:1.1,1.0", "heart-beat:10000,10000").build();
@@ -81,18 +107,7 @@ public class StompProtocolHandlerTests {
assertArrayEquals(new long[] {10000, 10000}, headers.getHeartbeat());
assertEquals(new HashSet<>(Arrays.asList("1.1","1.0")), headers.getAcceptVersion());
// Check CONNECTED reply
assertEquals(1, this.session.getSentMessages().size());
textMessage = (TextMessage) this.session.getSentMessages().get(0);
Message<?> message = new StompDecoder().decode(ByteBuffer.wrap(textMessage.getPayload().getBytes()));
StompHeaderAccessor replyHeaders = StompHeaderAccessor.wrap(message);
assertEquals(StompCommand.CONNECTED, replyHeaders.getCommand());
assertEquals("1.1", replyHeaders.getVersion());
assertArrayEquals(new long[] {0, 0}, replyHeaders.getHeartbeat());
assertEquals("joe", replyHeaders.getNativeHeader("user-name").get(0));
assertEquals("s1", replyHeaders.getNativeHeader("queue-suffix").get(0));
assertEquals(0, this.session.getSentMessages().size());
}
}