GH-7925: Make message history header as mutable

Fixes: #7925

The `MessageHistory.write()` creates not only a new instance of the `MessageHistory`,
but also a new copy of the whole message.
This significantly impacts the performance when we have too many components to track

* Make `MessageHistory` as append-only container and create a new instance (plus message)
only on the first track when no prior history is present
* Change `WireTap`, `BroadcastingDispatcher`, `AbstractMessageRouter` and `AbstractMessageSplitter`
to use a new `AbstractIntegrationMessageBuilder.cloneMessageHistoryIfAny()` API for every branch a message
is produced.
Essentially, create a new message with copy of the message history to let that downstream sub-flow
have its own trace
* Modify failed unit tests for a new logic where message history header is not immutable anymore
* This also fixes an `AbstractMessageSplitter` for propagating its track into messages it emits

* * Do not clone message history header if only one consume in multi-publish
* Fix typos in docs
This commit is contained in:
Artem Bilan
2024-03-05 15:05:34 -05:00
committed by GitHub
parent 3fbe917e6c
commit 2731e9411e
12 changed files with 214 additions and 134 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2020 the original author or authors.
* Copyright 2002-2024 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.
@@ -23,8 +23,11 @@ import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.integration.core.MessageSelector;
import org.springframework.integration.history.MessageHistory;
import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.support.channel.ChannelResolverUtils;
import org.springframework.integration.support.management.ManageableLifecycle;
import org.springframework.integration.support.utils.IntegrationUtils;
import org.springframework.jmx.export.annotation.ManagedAttribute;
import org.springframework.jmx.export.annotation.ManagedOperation;
import org.springframework.jmx.export.annotation.ManagedResource;
@@ -57,6 +60,8 @@ public class WireTap implements ChannelInterceptor, ManageableLifecycle, VetoCap
private BeanFactory beanFactory;
private MessageBuilderFactory messageBuilderFactory;
private volatile boolean running = true;
@@ -162,9 +167,19 @@ public class WireTap implements ChannelInterceptor, ManageableLifecycle, VetoCap
return message;
}
if (this.running && (this.selector == null || this.selector.accept(message))) {
boolean sent = (this.timeout >= 0)
? wireTapChannel.send(message, this.timeout)
: wireTapChannel.send(message);
Message<?> messageToSend = message;
if (message.getHeaders().containsKey(MessageHistory.HEADER_NAME)) {
messageToSend =
getMessageBuilderFactory()
.fromMessage(message)
.cloneMessageHistoryIfAny()
.build();
}
boolean sent =
(this.timeout >= 0)
? wireTapChannel.send(messageToSend, this.timeout)
: wireTapChannel.send(messageToSend);
if (!sent && LOGGER.isWarnEnabled()) {
LOGGER.warn("failed to send message to WireTap channel '" + wireTapChannel + "'");
}
@@ -174,7 +189,6 @@ public class WireTap implements ChannelInterceptor, ManageableLifecycle, VetoCap
@Override
public boolean shouldIntercept(String beanName, InterceptableChannel channel) {
return !getChannel().equals(channel);
}
@@ -189,4 +203,11 @@ public class WireTap implements ChannelInterceptor, ManageableLifecycle, VetoCap
return this.channel;
}
private MessageBuilderFactory getMessageBuilderFactory() {
if (this.messageBuilderFactory == null) {
this.messageBuilderFactory = IntegrationUtils.getMessageBuilderFactory(this.beanFactory);
}
return this.messageBuilderFactory;
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2023 the original author or authors.
* Copyright 2002-2024 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.
@@ -17,18 +17,23 @@
package org.springframework.integration.dispatcher;
import java.util.Collection;
import java.util.UUID;
import java.util.concurrent.Executor;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.integration.MessageDispatchingException;
import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.integration.history.MessageHistory;
import org.springframework.integration.support.AbstractIntegrationMessageBuilder;
import org.springframework.integration.support.DefaultMessageBuilderFactory;
import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.support.MessageDecorator;
import org.springframework.integration.support.utils.IntegrationUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.support.MessageHandlingRunnable;
import org.springframework.util.Assert;
@@ -153,13 +158,28 @@ public class BroadcastingDispatcher extends AbstractDispatcher implements BeanFa
}
int sequenceSize = handlers.size();
Message<?> messageToSend = message;
MessageHeaders messageHeaders = message.getHeaders();
UUID correlationKey = messageHeaders.getId();
boolean hasMessageHistory = messageHeaders.containsKey(MessageHistory.HEADER_NAME) && sequenceSize > 1;
for (MessageHandler handler : handlers) {
if (this.applySequence) {
messageToSend = getMessageBuilderFactory()
.fromMessage(message)
.pushSequenceDetails(message.getHeaders().getId(), sequenceNumber++, sequenceSize)
.build();
if (this.applySequence || hasMessageHistory) {
AbstractIntegrationMessageBuilder<?> builder =
getMessageBuilderFactory()
.fromMessage(message);
if (this.applySequence) {
builder.pushSequenceDetails(
correlationKey == null ? IntegrationObjectSupport.generateId() : correlationKey,
sequenceNumber++, sequenceSize);
}
if (hasMessageHistory) {
builder.cloneMessageHistoryIfAny();
}
messageToSend = builder.build();
}
if (message instanceof MessageDecorator messageDecorator) {
messageToSend = messageDecorator.decorateMessage(messageToSend);
}

View File

@@ -20,6 +20,7 @@ import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -465,6 +466,13 @@ public abstract class AbstractMessageProducingHandler extends AbstractMessageHan
}
else {
builder = getMessageBuilderFactory().withPayload(output);
// Assuming that message in the payload collection is a copy of request message.
if (output instanceof Iterable<?> iterable) {
Iterator<?> iterator = iterable.iterator();
if (iterator.hasNext() && iterator.next() instanceof Message<?>) {
builder = builder.cloneMessageHistoryIfAny();
}
}
}
if (!this.noHeadersPropagation &&
(shouldCopyRequestHeaders() ||

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2022 the original author or authors.
* Copyright 2002-2024 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.
@@ -16,6 +16,7 @@
package org.springframework.integration.history;
import java.io.Serial;
import java.io.Serializable;
import java.util.ArrayList;
import java.util.Collection;
@@ -52,8 +53,9 @@ import org.springframework.util.Assert;
*
* @since 2.0
*/
public final class MessageHistory implements List<Properties>, Serializable {
public final class MessageHistory implements List<Properties>, Serializable, Cloneable {
@Serial
private static final long serialVersionUID = -2340400235574314134L;
private static final Log LOGGER = LogFactory.getLog(MessageHistory.class);
@@ -94,51 +96,52 @@ public final class MessageHistory implements List<Properties>, Serializable {
Assert.notNull(component, "Component must not be null");
Properties metadata = extractMetadata(component);
if (!metadata.isEmpty()) {
MessageHistory previousHistory = message.getHeaders().get(HEADER_NAME, MessageHistory.class);
List<Properties> components =
previousHistory != null
? new ArrayList<>(previousHistory)
: new ArrayList<>();
components.add(metadata);
MessageHistory history = new MessageHistory(components);
if (message instanceof MutableMessage) {
message.getHeaders().put(HEADER_NAME, history);
}
else if (message instanceof ErrorMessage) {
ErrorMessage errorMessage = (ErrorMessage) message;
IntegrationMessageHeaderAccessor headerAccessor = new IntegrationMessageHeaderAccessor(message);
headerAccessor.setHeader(HEADER_NAME, history);
Throwable payload = errorMessage.getPayload();
Message<?> originalMessage = errorMessage.getOriginalMessage();
if (originalMessage != null) {
errorMessage = new ErrorMessage(payload, headerAccessor.toMessageHeaders(), originalMessage);
}
else {
errorMessage = new ErrorMessage(payload, headerAccessor.toMessageHeaders());
}
message = (Message<T>) errorMessage;
}
else if (message instanceof AdviceMessage) {
IntegrationMessageHeaderAccessor headerAccessor = new IntegrationMessageHeaderAccessor(message);
headerAccessor.setHeader(HEADER_NAME, history);
message = new AdviceMessage<T>(message.getPayload(), headerAccessor.toMessageHeaders(),
((AdviceMessage<?>) message).getInputMessage());
MessageHistory messageHistory = message.getHeaders().get(HEADER_NAME, MessageHistory.class);
if (messageHistory != null) {
messageHistory.components.add(metadata);
}
else {
if (!(message instanceof GenericMessage) &&
(messageBuilderFactory instanceof DefaultMessageBuilderFactory ||
messageBuilderFactory instanceof MutableMessageBuilderFactory)
&& LOGGER.isWarnEnabled()) {
List<Properties> components = new ArrayList<>();
components.add(metadata);
messageHistory = new MessageHistory(components);
LOGGER.warn("MessageHistory rebuilds the message and produces the result of the [" +
messageBuilderFactory + "], not an instance of the provided type [" +
message.getClass() + "]. Consider to supply a custom MessageBuilderFactory " +
"to retain custom messages during MessageHistory tracking.");
if (message instanceof MutableMessage) {
message.getHeaders().put(HEADER_NAME, messageHistory);
}
else if (message instanceof ErrorMessage errorMessage) {
IntegrationMessageHeaderAccessor headerAccessor = new IntegrationMessageHeaderAccessor(message);
headerAccessor.setHeader(HEADER_NAME, messageHistory);
Throwable payload = errorMessage.getPayload();
Message<?> originalMessage = errorMessage.getOriginalMessage();
if (originalMessage != null) {
errorMessage = new ErrorMessage(payload, headerAccessor.toMessageHeaders(), originalMessage);
}
else {
errorMessage = new ErrorMessage(payload, headerAccessor.toMessageHeaders());
}
message = (Message<T>) errorMessage;
}
else if (message instanceof AdviceMessage<?> adviceMessage) {
IntegrationMessageHeaderAccessor headerAccessor = new IntegrationMessageHeaderAccessor(message);
headerAccessor.setHeader(HEADER_NAME, messageHistory);
message = new AdviceMessage<>(message.getPayload(), headerAccessor.toMessageHeaders(),
adviceMessage.getInputMessage());
}
else {
if (!(message instanceof GenericMessage) &&
(messageBuilderFactory instanceof DefaultMessageBuilderFactory ||
messageBuilderFactory instanceof MutableMessageBuilderFactory)
&& LOGGER.isWarnEnabled()) {
LOGGER.warn("MessageHistory rebuilds the message and produces the result of the [" +
messageBuilderFactory + "], not an instance of the provided type [" +
message.getClass() + "]. Consider to supply a custom MessageBuilderFactory " +
"to retain custom messages during MessageHistory tracking.");
}
message = messageBuilderFactory.fromMessage(message)
.setHeader(HEADER_NAME, messageHistory)
.build();
}
message = messageBuilderFactory.fromMessage(message)
.setHeader(HEADER_NAME, history)
.build();
}
}
return message;
@@ -216,15 +219,19 @@ public final class MessageHistory implements List<Properties>, Serializable {
return this.components.lastIndexOf(o);
}
@Override
public Object clone() {
return new MessageHistory(new ArrayList<>(this.components));
}
@Override
public boolean equals(Object o) {
if (this == o) {
return true;
}
if (!(o instanceof MessageHistory)) {
if (!(o instanceof MessageHistory that)) {
return false;
}
MessageHistory that = (MessageHistory) o;
return this.components.equals(that.components);
}
@@ -318,6 +325,7 @@ public final class MessageHistory implements List<Properties>, Serializable {
*/
public static class Entry extends Properties {
@Serial
private static final long serialVersionUID = -8225834391885601079L;
public String getName() {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2023 the original author or authors.
* Copyright 2002-2024 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.
@@ -27,6 +27,8 @@ import org.springframework.core.convert.support.DefaultConversionService;
import org.springframework.integration.IntegrationPatternType;
import org.springframework.integration.core.MessagingTemplate;
import org.springframework.integration.handler.AbstractMessageHandler;
import org.springframework.integration.history.MessageHistory;
import org.springframework.integration.support.AbstractIntegrationMessageBuilder;
import org.springframework.integration.support.management.IntegrationManagedResource;
import org.springframework.jmx.export.annotation.ManagedResource;
import org.springframework.messaging.Message;
@@ -193,19 +195,27 @@ public abstract class AbstractMessageRouter extends AbstractMessageHandler imple
if (results != null) {
int sequenceSize = results.size();
int sequenceNumber = 1;
Message<?> messageToSend = message;
UUID correlationKey = message.getHeaders().getId();
boolean hasMessageHistory = message.getHeaders().containsKey(MessageHistory.HEADER_NAME) && sequenceSize > 1;
for (MessageChannel channel : results) {
final Message<?> messageToSend;
if (!this.applySequence) {
messageToSend = message;
}
else {
UUID id = message.getHeaders().getId();
messageToSend = getMessageBuilderFactory()
.fromMessage(message)
.pushSequenceDetails(id == null ? generateId() : id,
sequenceNumber++, sequenceSize)
.build();
if (this.applySequence || hasMessageHistory) {
AbstractIntegrationMessageBuilder<?> builder =
getMessageBuilderFactory()
.fromMessage(message);
if (this.applySequence) {
builder.pushSequenceDetails(correlationKey == null ? generateId() : correlationKey,
sequenceNumber++, sequenceSize);
}
if (hasMessageHistory) {
builder.cloneMessageHistoryIfAny();
}
messageToSend = builder.build();
}
if (channel != null) {
sent |= doSend(channel, messageToSend);
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2023 the original author or authors.
* Copyright 2002-2024 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.
@@ -35,6 +35,7 @@ import org.springframework.integration.IntegrationPatternType;
import org.springframework.integration.channel.ReactiveStreamsSubscribableChannel;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
import org.springframework.integration.handler.DiscardingMessageHandler;
import org.springframework.integration.history.MessageHistory;
import org.springframework.integration.support.AbstractIntegrationMessageBuilder;
import org.springframework.integration.support.json.JacksonPresent;
import org.springframework.integration.util.FunctionIterator;
@@ -276,10 +277,14 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
Object correlationId, int sequenceNumber, int sequenceSize) {
AbstractIntegrationMessageBuilder<?> builder = messageBuilderForReply(item);
builder.copyHeadersIfAbsent(headers);
builder.setHeader(MessageHistory.HEADER_NAME, headers.get(MessageHistory.HEADER_NAME))
.copyHeadersIfAbsent(headers)
.cloneMessageHistoryIfAny();
if (this.applySequence) {
builder.pushSequenceDetails(correlationId, sequenceNumber, sequenceSize);
}
return builder;
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2014-2021 the original author or authors.
* Copyright 2014-2024 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.
@@ -25,6 +25,7 @@ import java.util.List;
import java.util.Map;
import org.springframework.integration.IntegrationMessageHeaderAccessor;
import org.springframework.integration.history.MessageHistory;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
@@ -165,6 +166,22 @@ public abstract class AbstractIntegrationMessageBuilder<T> {
return copyHeadersIfAbsent(headers);
}
/**
* Make a copy of {@link MessageHistory} header (if present) for a new message to build.
* @return the current {@link AbstractIntegrationMessageBuilder}.
* @since 6.3
*/
public AbstractIntegrationMessageBuilder<T> cloneMessageHistoryIfAny() {
MessageHistory messageHistory = getHeader(MessageHistory.HEADER_NAME, MessageHistory.class);
if (messageHistory != null) {
return removeHeader(MessageHistory.HEADER_NAME)
.setHeader(MessageHistory.HEADER_NAME, messageHistory.clone());
}
return this;
}
@Nullable
protected abstract List<List<Object>> getSequenceDetails();

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2014-2022 the original author or authors.
* Copyright 2014-2024 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.
@@ -186,10 +186,8 @@ public class MessagingAnnotationsWithBeanAnnotationTests {
assertThat(messageHistory).isNotNull();
String messageHistoryString = messageHistory.toString();
assertThat(messageHistoryString)
.contains("routerChannel", "filterChannel", "aggregatorChannel", "serviceChannel")
.doesNotContain("discardChannel")
// history header is not overridden in splitter for individual message from message group emitted before
.doesNotContain("splitterChannel");
.contains("routerChannel", "filterChannel", "aggregatorChannel", "splitterChannel", "serviceChannel")
.doesNotContain("discardChannel");
}
assertThat(this.skippedServiceActivator).isNull();

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2022 the original author or authors.
* Copyright 2002-2024 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.
@@ -96,7 +96,7 @@ public class MessageHistoryTests {
Message<Throwable> result2 = MessageHistory.write(result1, new TestComponent(2));
assertThat(result2).isInstanceOf(ErrorMessage.class);
assertThat(result2).isNotSameAs(original);
assertThat(result2).isNotSameAs(result1);
assertThat(result2).isSameAs(result1);
assertThat(result2.getPayload()).isSameAs(original.getPayload());
assertThat(result1).extracting("originalMessage").isSameAs(originalMessage);
MessageHistory history2 = MessageHistory.read(result2);
@@ -122,20 +122,14 @@ public class MessageHistoryTests {
assertThat(result2).isNotSameAs(original);
assertThat(result2.getPayload()).isSameAs(original.getPayload());
assertThat(((AdviceMessage<?>) result2).getInputMessage()).isSameAs(original.getInputMessage());
assertThat(result2).isNotSameAs(result1);
assertThat(result2).isSameAs(result1);
MessageHistory history2 = MessageHistory.read(result2);
assertThat(history2).isNotNull();
assertThat(history2.toString()).isEqualTo("testComponent-1,testComponent-2");
}
private static class TestComponent implements NamedComponent {
private final int id;
TestComponent(int id) {
this.id = id;
}
private record TestComponent(int id) implements NamedComponent {
@Override
public String getComponentName() {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2022 the original author or authors.
* Copyright 2002-2024 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.
@@ -21,8 +21,8 @@ import java.util.Collections;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.Before;
import org.junit.Test;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
import org.springframework.core.task.TaskExecutor;
@@ -35,10 +35,12 @@ import org.springframework.messaging.MessagingException;
import org.springframework.messaging.support.GenericMessage;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.fail;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
/**
@@ -52,26 +54,27 @@ public class BroadcastingDispatcherTests {
private BroadcastingDispatcher dispatcher;
private final TaskExecutor taskExecutorMock = Mockito.mock(TaskExecutor.class);
private final TaskExecutor taskExecutorMock = Mockito.mock();
private final Message<?> messageMock = Mockito.mock(Message.class);
private final Message<?> messageMock = Mockito.mock();
private final MessageHandler targetMock1 = Mockito.mock(MessageHandler.class);
private final MessageHandler targetMock1 = Mockito.mock();
private final MessageHandler targetMock2 = Mockito.mock(MessageHandler.class);
private final MessageHandler targetMock2 = Mockito.mock();
private final MessageHandler targetMock3 = Mockito.mock(MessageHandler.class);
private final MessageHandler targetMock3 = Mockito.mock();
@Before
@BeforeEach
public void init() {
Mockito.reset(taskExecutorMock, messageMock, taskExecutorMock, targetMock1, targetMock2, targetMock3);
given(messageMock.getHeaders()).willReturn(mock());
defaultTaskExecutorMock();
}
@Test
public void singleTargetWithoutTaskExecutor() throws Exception {
public void singleTargetWithoutTaskExecutor() {
dispatcher = new BroadcastingDispatcher();
dispatcher.addHandler(targetMock1);
dispatcher.dispatch(messageMock);
@@ -79,7 +82,7 @@ public class BroadcastingDispatcherTests {
}
@Test
public void singleTargetWithTaskExecutor() throws Exception {
public void singleTargetWithTaskExecutor() {
dispatcher = new BroadcastingDispatcher(taskExecutorMock);
dispatcher.addHandler(targetMock1);
dispatcher.dispatch(messageMock);
@@ -202,12 +205,12 @@ public class BroadcastingDispatcherTests {
@Test
public void applySequenceDisabledByDefault() {
BroadcastingDispatcher dispatcher = new BroadcastingDispatcher();
final List<Message<?>> messages = Collections.synchronizedList(new ArrayList<Message<?>>());
final List<Message<?>> messages = Collections.synchronizedList(new ArrayList<>());
MessageHandler target1 = new MessageStoringTestEndpoint(messages);
MessageHandler target2 = new MessageStoringTestEndpoint(messages);
dispatcher.addHandler(target1);
dispatcher.addHandler(target2);
dispatcher.dispatch(new GenericMessage<String>("test"));
dispatcher.dispatch(new GenericMessage<>("test"));
assertThat(messages.size()).isEqualTo(2);
assertThat(new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceNumber()).isEqualTo(0);
assertThat(new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceSize()).isEqualTo(0);
@@ -219,14 +222,14 @@ public class BroadcastingDispatcherTests {
public void applySequenceEnabled() {
BroadcastingDispatcher dispatcher = new BroadcastingDispatcher();
dispatcher.setApplySequence(true);
final List<Message<?>> messages = Collections.synchronizedList(new ArrayList<Message<?>>());
final List<Message<?>> messages = Collections.synchronizedList(new ArrayList<>());
MessageHandler target1 = new MessageStoringTestEndpoint(messages);
MessageHandler target2 = new MessageStoringTestEndpoint(messages);
MessageHandler target3 = new MessageStoringTestEndpoint(messages);
dispatcher.addHandler(target1);
dispatcher.addHandler(target2);
dispatcher.addHandler(target3);
Message<?> inputMessage = new GenericMessage<String>("test");
Message<?> inputMessage = new GenericMessage<>("test");
Object originalId = inputMessage.getHeaders().getId();
dispatcher.dispatch(inputMessage);
assertThat(messages.size()).isEqualTo(3);
@@ -251,13 +254,11 @@ public class BroadcastingDispatcherTests {
dispatcher.addHandler(targetMock1);
doThrow(new MessagingException("Mock Exception"))
.when(targetMock1).handleMessage(eq(messageMock));
try {
dispatcher.dispatch(messageMock);
fail("Expected Exception");
}
catch (MessagingException e) {
assertThat(e.getFailedMessage()).isEqualTo(messageMock);
}
assertThatExceptionOfType(MessagingException.class)
.isThrownBy(() -> dispatcher.dispatch(messageMock))
.extracting(MessagingException::getFailedMessage)
.isEqualTo(messageMock);
}
/**
@@ -272,13 +273,11 @@ public class BroadcastingDispatcherTests {
Message<String> dontReplaceThisMessage = MessageBuilder.withPayload("x").build();
doThrow(new MessagingException(dontReplaceThisMessage, "Mock Exception"))
.when(targetMock1).handleMessage(eq(messageMock));
try {
dispatcher.dispatch(messageMock);
fail("Expected Exception");
}
catch (MessagingException e) {
assertThat(e.getFailedMessage()).isEqualTo(dontReplaceThisMessage);
}
assertThatExceptionOfType(MessagingException.class)
.isThrownBy(() -> dispatcher.dispatch(messageMock))
.extracting(MessagingException::getFailedMessage)
.isEqualTo(dontReplaceThisMessage);
}
/**
@@ -306,13 +305,10 @@ public class BroadcastingDispatcherTests {
@Test
public void testNoHandlerWithRequiredSubscriber() {
dispatcher = new BroadcastingDispatcher(true);
try {
dispatcher.dispatch(messageMock);
fail("Expected Exception");
}
catch (MessageDispatchingException exception) {
assertThat(exception.getFailedMessage()).isEqualTo(messageMock);
}
assertThatExceptionOfType(MessageDispatchingException.class)
.isThrownBy(() -> dispatcher.dispatch(messageMock))
.extracting(MessagingException::getFailedMessage)
.isEqualTo(messageMock);
}
/**
@@ -322,13 +318,10 @@ public class BroadcastingDispatcherTests {
@Test
public void testNoHandlerWithExecutorWithRequiredSubscriber() {
dispatcher = new BroadcastingDispatcher(taskExecutorMock, true);
try {
dispatcher.dispatch(messageMock);
fail("Expected Exception");
}
catch (MessageDispatchingException exception) {
assertThat(exception.getFailedMessage()).isEqualTo(messageMock);
}
assertThatExceptionOfType(MessageDispatchingException.class)
.isThrownBy(() -> dispatcher.dispatch(messageMock))
.extracting(MessagingException::getFailedMessage)
.isEqualTo(messageMock);
}
private void defaultTaskExecutorMock() {