GH-2388: Fix HeaderEnricherSpec for adviceChain
Fixes spring-projects/spring-integration#2388 During `HeaderEnricherSpec` refactoring the `adviceChain` population has been missed, alongside with many other `AbstractReplyProducingMessageHandler` options from the `ConsumerEndpointSpec` * Refactor `HeaderEnricherSpec` to build `HeaderEnricher` and an appropriate `MessageTransformingHandler` from the ctor to be able to pick up an `adviceChain` automatically in the `ConsumerEndpointSpec.get()` **Cherry-pick to 5.0.x**
This commit is contained in:
committed by
Gary Russell
parent
4552a80095
commit
a4eb4ec9b1
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2017 the original author or authors.
|
||||
* Copyright 2016-2018 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.
|
||||
@@ -56,10 +56,6 @@ public abstract class ConsumerEndpointSpec<S extends ConsumerEndpointSpec<S, H>,
|
||||
|
||||
protected ConsumerEndpointSpec(H messageHandler) {
|
||||
super(messageHandler);
|
||||
this.endpointFactoryBean.setAdviceChain(this.adviceChain);
|
||||
if (messageHandler instanceof AbstractReplyProducingMessageHandler) {
|
||||
((AbstractReplyProducingMessageHandler) messageHandler).setAdviceChain(this.adviceChain);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -270,6 +266,10 @@ public abstract class ConsumerEndpointSpec<S extends ConsumerEndpointSpec<S, H>,
|
||||
|
||||
@Override
|
||||
protected Tuple2<ConsumerEndpointFactoryBean, H> doGet() {
|
||||
this.endpointFactoryBean.setAdviceChain(this.adviceChain);
|
||||
if (this.handler instanceof AbstractReplyProducingMessageHandler) {
|
||||
((AbstractReplyProducingMessageHandler) this.handler).setAdviceChain(this.adviceChain);
|
||||
}
|
||||
this.endpointFactoryBean.setHandler(this.handler);
|
||||
return super.doGet();
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2017 the original author or authors.
|
||||
* Copyright 2016-2018 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.
|
||||
@@ -55,14 +55,11 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
|
||||
private final Map<String, HeaderValueMessageProcessor<?>> headerToAdd = new HashMap<>();
|
||||
|
||||
private boolean defaultOverwrite = false;
|
||||
|
||||
private boolean shouldSkipNulls = true;
|
||||
|
||||
private MessageProcessor<?> messageProcessor;
|
||||
private final HeaderEnricher headerEnricher = new HeaderEnricher(this.headerToAdd);
|
||||
|
||||
HeaderEnricherSpec() {
|
||||
super(null);
|
||||
this.handler = new MessageTransformingHandler(this.headerEnricher);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -73,7 +70,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
* @see HeaderEnricher#setDefaultOverwrite(boolean)
|
||||
*/
|
||||
public HeaderEnricherSpec defaultOverwrite(boolean defaultOverwrite) {
|
||||
this.defaultOverwrite = defaultOverwrite;
|
||||
this.headerEnricher.setDefaultOverwrite(defaultOverwrite);
|
||||
return _this();
|
||||
}
|
||||
|
||||
@@ -83,7 +80,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
* @see HeaderEnricher#setShouldSkipNulls(boolean)
|
||||
*/
|
||||
public HeaderEnricherSpec shouldSkipNulls(boolean shouldSkipNulls) {
|
||||
this.shouldSkipNulls = shouldSkipNulls;
|
||||
this.headerEnricher.setShouldSkipNulls(shouldSkipNulls);
|
||||
return _this();
|
||||
}
|
||||
|
||||
@@ -97,7 +94,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
* @see HeaderEnricher#setMessageProcessor(MessageProcessor)
|
||||
*/
|
||||
public HeaderEnricherSpec messageProcessor(MessageProcessor<?> messageProcessor) {
|
||||
this.messageProcessor = messageProcessor;
|
||||
this.headerEnricher.setMessageProcessor(messageProcessor);
|
||||
return _this();
|
||||
}
|
||||
|
||||
@@ -110,7 +107,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
* @see #messageProcessor(MessageProcessor)
|
||||
*/
|
||||
public HeaderEnricherSpec messageProcessor(String expression) {
|
||||
return messageProcessor(new ExpressionEvaluatingMessageProcessor<>(PARSER.parseExpression(expression)));
|
||||
return messageProcessor(new ExpressionEvaluatingMessageProcessor<>(expression));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -371,6 +368,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
*/
|
||||
public <P> HeaderEnricherSpec headerFunction(String name, Function<Message<P>, Object> function,
|
||||
Boolean overwrite) {
|
||||
|
||||
return headerExpression(name, new FunctionExpression<>(function), overwrite);
|
||||
}
|
||||
|
||||
@@ -391,6 +389,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
*/
|
||||
public <V> HeaderEnricherSpec header(String headerName,
|
||||
HeaderValueMessageProcessor<V> headerValueMessageProcessor) {
|
||||
|
||||
Assert.hasText(headerName, "'headerName' must not be empty");
|
||||
this.headerToAdd.put(headerName, headerValueMessageProcessor);
|
||||
return _this();
|
||||
@@ -432,14 +431,7 @@ public class HeaderEnricherSpec extends ConsumerEndpointSpec<HeaderEnricherSpec,
|
||||
|
||||
@Override
|
||||
protected Tuple2<ConsumerEndpointFactoryBean, MessageTransformingHandler> doGet() {
|
||||
HeaderEnricher headerEnricher = new HeaderEnricher(new HashMap<>(this.headerToAdd));
|
||||
headerEnricher.setDefaultOverwrite(this.defaultOverwrite);
|
||||
headerEnricher.setShouldSkipNulls(this.shouldSkipNulls);
|
||||
headerEnricher.setMessageProcessor(this.messageProcessor);
|
||||
|
||||
this.componentsToRegister.put(headerEnricher, null);
|
||||
|
||||
this.handler = new MessageTransformingHandler(headerEnricher);
|
||||
this.componentsToRegister.put(this.headerEnricher, null);
|
||||
return super.doGet();
|
||||
}
|
||||
|
||||
|
||||
@@ -51,6 +51,7 @@ import org.springframework.integration.dsl.IntegrationFlow;
|
||||
import org.springframework.integration.dsl.IntegrationFlows;
|
||||
import org.springframework.integration.dsl.Transformers;
|
||||
import org.springframework.integration.dsl.channel.MessageChannels;
|
||||
import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvice;
|
||||
import org.springframework.integration.handler.advice.IdempotentReceiverInterceptor;
|
||||
import org.springframework.integration.selector.MetadataStoreSelector;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
@@ -91,7 +92,6 @@ public class TransformerTests {
|
||||
private PollableChannel enricherErrorChannel;
|
||||
|
||||
|
||||
|
||||
@Test
|
||||
public void testContentEnricher() {
|
||||
QueueChannel replyChannel = new QueueChannel();
|
||||
@@ -218,6 +218,9 @@ public class TransformerTests {
|
||||
@Autowired
|
||||
private PollableChannel idempotentDiscardChannel;
|
||||
|
||||
@Autowired
|
||||
private PollableChannel adviceChannel;
|
||||
|
||||
@Test
|
||||
public void transformWithHeader() {
|
||||
QueueChannel replyChannel = new QueueChannel();
|
||||
@@ -240,6 +243,7 @@ public class TransformerTests {
|
||||
}
|
||||
|
||||
assertNotNull(this.idempotentDiscardChannel.receive(10000));
|
||||
assertNotNull(this.adviceChannel.receive(10000));
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@@ -328,7 +332,7 @@ public class TransformerTests {
|
||||
return f -> f
|
||||
.enrichHeaders(h -> h
|
||||
.header("Foo", "Bar")
|
||||
.advice(idempotentReceiverInterceptor()))
|
||||
.advice(idempotentReceiverInterceptor(), requestHandlerAdvice()))
|
||||
.transform(new PojoTransformer());
|
||||
}
|
||||
|
||||
@@ -346,6 +350,26 @@ public class TransformerTests {
|
||||
return idempotentReceiverInterceptor;
|
||||
}
|
||||
|
||||
@Bean
|
||||
public AbstractRequestHandlerAdvice requestHandlerAdvice() {
|
||||
return new AbstractRequestHandlerAdvice() {
|
||||
|
||||
@Override
|
||||
protected Object doInvoke(ExecutionCallback callback, Object target, Message<?> message)
|
||||
throws Exception {
|
||||
|
||||
adviceChannel().send(message);
|
||||
return callback.execute();
|
||||
}
|
||||
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public PollableChannel adviceChannel() {
|
||||
return new QueueChannel();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public PollableChannel subFlowTestReplyChannel() {
|
||||
return new QueueChannel();
|
||||
|
||||
Reference in New Issue
Block a user