OrderedMessageChannelDecorator doesn't preclude send limits
Closes gh-25581
This commit is contained in:
@@ -142,7 +142,7 @@ public abstract class AbstractBrokerMessageHandler
|
||||
* @since 5.1
|
||||
*/
|
||||
public void setPreservePublishOrder(boolean preservePublishOrder) {
|
||||
OrderedMessageSender.configureOutboundChannel(this.clientOutboundChannel, preservePublishOrder);
|
||||
OrderedMessageChannelDecorator.configureInterceptor(this.clientOutboundChannel, preservePublishOrder);
|
||||
this.preservePublishOrder = preservePublishOrder;
|
||||
}
|
||||
|
||||
@@ -298,7 +298,7 @@ public abstract class AbstractBrokerMessageHandler
|
||||
*/
|
||||
protected MessageChannel getClientOutboundChannelForSession(String sessionId) {
|
||||
return this.preservePublishOrder ?
|
||||
new OrderedMessageSender(getClientOutboundChannel(), logger) : getClientOutboundChannel();
|
||||
new OrderedMessageChannelDecorator(getClientOutboundChannel(), logger) : getClientOutboundChannel();
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2018 the original author or authors.
|
||||
* Copyright 2002-2020 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.
|
||||
@@ -33,15 +33,17 @@ import org.springframework.messaging.support.MessageHeaderAccessor;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* Submit messages to an {@link ExecutorSubscribableChannel}, one at a time.
|
||||
* The channel must have been configured with {@link #configureOutboundChannel}.
|
||||
* Decorator for an {@link ExecutorSubscribableChannel} that ensures messages
|
||||
* are processed in the order they were published to the channel. Messages are
|
||||
* sent one at a time with the next one released when the prevoius has been
|
||||
* processed. This decorator is intended to be applied per session.
|
||||
*
|
||||
* @author Rossen Stoyanchev
|
||||
* @since 5.1
|
||||
*/
|
||||
class OrderedMessageSender implements MessageChannel {
|
||||
public class OrderedMessageChannelDecorator implements MessageChannel {
|
||||
|
||||
static final String COMPLETION_TASK_HEADER = "simpSendCompletionTask";
|
||||
private static final String NEXT_MESSAGE_TASK_HEADER = "simpNextMessageTask";
|
||||
|
||||
|
||||
private final MessageChannel channel;
|
||||
@@ -53,7 +55,7 @@ class OrderedMessageSender implements MessageChannel {
|
||||
private final AtomicBoolean sendInProgress = new AtomicBoolean(false);
|
||||
|
||||
|
||||
public OrderedMessageSender(MessageChannel channel, Log logger) {
|
||||
public OrderedMessageChannelDecorator(MessageChannel channel, Log logger) {
|
||||
this.channel = channel;
|
||||
this.logger = logger;
|
||||
}
|
||||
@@ -84,10 +86,14 @@ class OrderedMessageSender implements MessageChannel {
|
||||
|
||||
private void sendNextMessage() {
|
||||
for (;;) {
|
||||
Message<?> message = this.messages.poll();
|
||||
Message<?> message = this.messages.peek();
|
||||
if (message != null) {
|
||||
try {
|
||||
addCompletionCallback(message);
|
||||
addNextMessageTaskHeader(message, () -> {
|
||||
if (removeMessage(message)) {
|
||||
sendNextMessage();
|
||||
}
|
||||
});
|
||||
if (this.channel.send(message)) {
|
||||
return;
|
||||
}
|
||||
@@ -97,9 +103,9 @@ class OrderedMessageSender implements MessageChannel {
|
||||
logger.error("Failed to send " + message, ex);
|
||||
}
|
||||
}
|
||||
removeMessage(message);
|
||||
}
|
||||
else {
|
||||
// We ran out of messages..
|
||||
this.sendInProgress.set(false);
|
||||
trySend();
|
||||
break;
|
||||
@@ -107,22 +113,40 @@ class OrderedMessageSender implements MessageChannel {
|
||||
}
|
||||
}
|
||||
|
||||
private void addCompletionCallback(Message<?> msg) {
|
||||
SimpMessageHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(msg, SimpMessageHeaderAccessor.class);
|
||||
Assert.isTrue(accessor != null && accessor.isMutable(), "Expected mutable SimpMessageHeaderAccessor");
|
||||
accessor.setHeader(COMPLETION_TASK_HEADER, (Runnable) this::sendNextMessage);
|
||||
private boolean removeMessage(Message<?> message) {
|
||||
Message<?> next = this.messages.peek();
|
||||
if (next == message) {
|
||||
this.messages.remove();
|
||||
return true;
|
||||
}
|
||||
else {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
private static void addNextMessageTaskHeader(Message<?> message, Runnable task) {
|
||||
SimpMessageHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, SimpMessageHeaderAccessor.class);
|
||||
Assert.isTrue(accessor != null && accessor.isMutable(), "Expected mutable SimpMessageHeaderAccessor");
|
||||
accessor.setHeader(NEXT_MESSAGE_TASK_HEADER, task);
|
||||
}
|
||||
|
||||
/**
|
||||
* Obtain the task to release the next message, if found.
|
||||
*/
|
||||
@Nullable
|
||||
public static Runnable getNextMessageTask(Message<?> message) {
|
||||
return (Runnable) message.getHeaders().get(OrderedMessageChannelDecorator.NEXT_MESSAGE_TASK_HEADER);
|
||||
}
|
||||
|
||||
/**
|
||||
* Install or remove an {@link ExecutorChannelInterceptor} that invokes a
|
||||
* completion task once the message is handled.
|
||||
* completion task, if found in the headers of the message.
|
||||
* @param channel the channel to configure
|
||||
* @param preservePublishOrder whether preserve order is on or off based on
|
||||
* which an interceptor is either added or removed.
|
||||
* @param preserveOrder whether preserve the order or publication; when
|
||||
* "true" an interceptor is inserted, when "false" it removed.
|
||||
*/
|
||||
static void configureOutboundChannel(MessageChannel channel, boolean preservePublishOrder) {
|
||||
if (preservePublishOrder) {
|
||||
public static void configureInterceptor(MessageChannel channel, boolean preserveOrder) {
|
||||
if (preserveOrder) {
|
||||
Assert.isInstanceOf(ExecutorSubscribableChannel.class, channel,
|
||||
"An ExecutorSubscribableChannel is required for `preservePublishOrder`");
|
||||
ExecutorSubscribableChannel execChannel = (ExecutorSubscribableChannel) channel;
|
||||
@@ -133,8 +157,7 @@ class OrderedMessageSender implements MessageChannel {
|
||||
else if (channel instanceof ExecutorSubscribableChannel) {
|
||||
ExecutorSubscribableChannel execChannel = (ExecutorSubscribableChannel) channel;
|
||||
execChannel.getInterceptors().stream().filter(i -> i instanceof CallbackInterceptor)
|
||||
.findFirst()
|
||||
.map(execChannel::removeInterceptor);
|
||||
.findFirst().map(execChannel::removeInterceptor);
|
||||
|
||||
}
|
||||
}
|
||||
@@ -144,9 +167,9 @@ class OrderedMessageSender implements MessageChannel {
|
||||
|
||||
@Override
|
||||
public void afterMessageHandled(
|
||||
Message<?> msg, MessageChannel ch, MessageHandler handler, @Nullable Exception ex) {
|
||||
Message<?> message, MessageChannel ch, MessageHandler handler, @Nullable Exception ex) {
|
||||
|
||||
Runnable task = (Runnable) msg.getHeaders().get(OrderedMessageSender.COMPLETION_TASK_HEADER);
|
||||
Runnable task = getNextMessageTask(message);
|
||||
if (task != null) {
|
||||
task.run();
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2019 the original author or authors.
|
||||
* Copyright 2002-2020 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.
|
||||
@@ -36,15 +36,16 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* Unit tests for {@link OrderedMessageSender}.
|
||||
* Unit tests for {@link OrderedMessageChannelDecorator}.
|
||||
* @author Rossen Stoyanchev
|
||||
* @see org.springframework.web.socket.messaging.OrderedMessageSendingIntegrationTests
|
||||
*/
|
||||
public class OrderedMessageSenderTests {
|
||||
public class OrderedMessageChannelDecoratorTests {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(OrderedMessageSenderTests.class);
|
||||
private static final Log logger = LogFactory.getLog(OrderedMessageChannelDecoratorTests.class);
|
||||
|
||||
|
||||
private OrderedMessageSender sender;
|
||||
private OrderedMessageChannelDecorator sender;
|
||||
|
||||
ExecutorSubscribableChannel channel = new ExecutorSubscribableChannel(this.executor);
|
||||
|
||||
@@ -59,9 +60,9 @@ public class OrderedMessageSenderTests {
|
||||
this.executor.afterPropertiesSet();
|
||||
|
||||
this.channel = new ExecutorSubscribableChannel(this.executor);
|
||||
OrderedMessageSender.configureOutboundChannel(this.channel, true);
|
||||
OrderedMessageChannelDecorator.configureInterceptor(this.channel, true);
|
||||
|
||||
this.sender = new OrderedMessageSender(this.channel, logger);
|
||||
this.sender = new OrderedMessageChannelDecorator(this.channel, logger);
|
||||
|
||||
}
|
||||
|
||||
@@ -89,9 +90,10 @@ public class OrderedMessageSenderTests {
|
||||
latch.countDown();
|
||||
return;
|
||||
}
|
||||
if (actual == 100 || actual == 200) {
|
||||
// Force messages to queue up periodically
|
||||
if (actual % 101 == 0) {
|
||||
try {
|
||||
Thread.sleep(200);
|
||||
Thread.sleep(50);
|
||||
}
|
||||
catch (InterruptedException ex) {
|
||||
result.set(ex.toString());
|
||||
Reference in New Issue
Block a user