INT-2763 Send AdviceMessage from Advice

The ExpressionEvaluatingMessageHandlerAdvice now sends an
AdviceMessage which has a payload of the original message, and
a property 'inputMessage' referencing the message that was
sent to the advised handler.

INT-2763 Expression Advice Reference Doc Update

Update to reflect the new behavior.
This commit is contained in:
Gary Russell
2012-09-21 17:55:50 +01:00
committed by Oleg Zhurakousky
parent c8da0532d7
commit 2340e6f442
5 changed files with 214 additions and 132 deletions

View File

@@ -22,6 +22,7 @@ import org.aopalliance.intercept.MethodInvocation;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.integration.Message;
import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.integration.core.MessageHandler;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
@@ -35,7 +36,8 @@ import org.springframework.integration.handler.AbstractReplyProducingMessageHand
* @since 2.2
*
*/
public abstract class AbstractRequestHandlerAdvice implements MethodInterceptor {
public abstract class AbstractRequestHandlerAdvice extends IntegrationObjectSupport
implements MethodInterceptor {
protected final Log logger = LogFactory.getLog(this.getClass());

View File

@@ -15,19 +15,21 @@
*/
package org.springframework.integration.handler.advice;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.context.expression.BeanFactoryResolver;
import org.springframework.expression.EvaluationContext;
import org.springframework.expression.Expression;
import org.springframework.expression.spel.standard.SpelExpressionParser;
import org.springframework.expression.spel.support.StandardEvaluationContext;
import org.springframework.integration.Message;
import org.springframework.integration.MessageChannel;
import org.springframework.integration.MessageHeaders;
import org.springframework.integration.MessagingException;
import org.springframework.integration.core.MessageHandler;
import org.springframework.integration.core.MessagingTemplate;
import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.util.StringUtils;
import org.springframework.integration.message.AdviceMessage;
import org.springframework.integration.message.ErrorMessage;
import org.springframework.integration.util.ExpressionUtils;
import org.springframework.util.Assert;
/**
* Used to advise {@link MessageHandler}s.
@@ -40,16 +42,15 @@ import org.springframework.util.StringUtils;
* @since 2.2
*
*/
public class ExpressionEvaluatingRequestHandlerAdvice extends AbstractRequestHandlerAdvice
implements BeanFactoryAware {
public class ExpressionEvaluatingRequestHandlerAdvice extends AbstractRequestHandlerAdvice {
private final ExpressionEvaluatingMessageProcessor<Object> onSuccessMessageProcessor;
private volatile Expression onSuccessExpression;
private final MessageChannel successChannel;
private volatile MessageChannel successChannel;
private final ExpressionEvaluatingMessageProcessor<Object> onFailureMessageProcessor;
private volatile Expression onFailureExpression;
private final MessageChannel failureChannel;
private volatile MessageChannel failureChannel;
private final MessagingTemplate messagingTemplate = new MessagingTemplate();
@@ -57,45 +58,28 @@ public class ExpressionEvaluatingRequestHandlerAdvice extends AbstractRequestHan
private volatile boolean returnFailureExpressionResult = false;
private volatile BeanFactory beanFactory;
private volatile boolean propagateOnSuccessEvaluationFailures;
/**
* @param onSuccessExpression
* @param successChannel
* @param onFailureExpression
* @param failureChannel
*/
public ExpressionEvaluatingRequestHandlerAdvice(Expression onSuccessExpression, MessageChannel successChannel,
Expression onFailureExpression, MessageChannel failureChannel) {
if (onSuccessExpression != null) {
this.onSuccessMessageProcessor = new ExpressionEvaluatingMessageProcessor<Object>(onSuccessExpression);
this.onSuccessMessageProcessor.setBeanFactory(this.beanFactory);
}
else {
this.onSuccessMessageProcessor = null;
}
this.successChannel = successChannel;
if (onFailureExpression != null) {
this.onFailureMessageProcessor = new ExpressionEvaluatingMessageProcessor<Object>(onFailureExpression);
onFailureMessageProcessor.setBeanFactory(this.beanFactory);
}
else {
this.onFailureMessageProcessor = null;
}
this.failureChannel = failureChannel;
private volatile EvaluationContext evaluationContext;
public void setOnSuccessExpression(String onSuccessExpression) {
Assert.notNull(onSuccessExpression, "'onSuccessExpression' must not be null");
this.onSuccessExpression = new SpelExpressionParser().parseExpression(onSuccessExpression);
}
public ExpressionEvaluatingRequestHandlerAdvice(String onSuccessExpression, MessageChannel successChannel,
String onFailureExpression, MessageChannel failureChannel) {
this(!StringUtils.hasText(onSuccessExpression) ? null :
new SpelExpressionParser().parseExpression(onSuccessExpression),
successChannel,
!StringUtils.hasText(onFailureExpression) ? null :
new SpelExpressionParser().parseExpression(onFailureExpression),
failureChannel);
public void setOnFailureExpression(String onFailureExpression) {
Assert.notNull(onFailureExpression, "'onFailureExpression' must not be null");
this.onFailureExpression = new SpelExpressionParser().parseExpression(onFailureExpression);
}
public void setSuccessChannel(MessageChannel successChannel) {
Assert.notNull(successChannel,"'successChannel' must not be null");
this.successChannel = successChannel;
}
public void setFailureChannel(MessageChannel failureChannel) {
Assert.notNull(failureChannel,"'failureChannel' must not be null");
this.failureChannel = failureChannel;
}
/**
@@ -126,22 +110,18 @@ public class ExpressionEvaluatingRequestHandlerAdvice extends AbstractRequestHan
this.propagateOnSuccessEvaluationFailures = propagateOnSuccessEvaluationFailures;
}
public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
this.beanFactory = beanFactory;
}
@Override
protected Object doInvoke(ExecutionCallback callback, Object target, Message<?> message) throws Exception {
try {
Object result = callback.execute();
if (this.onSuccessMessageProcessor != null) {
evaluateExpression(message, this.onSuccessMessageProcessor, this.successChannel, this.propagateOnSuccessEvaluationFailures);
if (this.onSuccessExpression != null) {
this.evaluateSuccessExpression(message);
}
return result;
}
catch (Exception e) {
if (this.onFailureMessageProcessor != null) {
Object evalResult = evaluateExpression(message, this.onFailureMessageProcessor, this.failureChannel, false);
if (this.onFailureExpression != null) {
Object evalResult = this.evaluateFailureExpression(message, e);
if (this.returnFailureExpressionResult) {
return evalResult;
}
@@ -153,28 +133,90 @@ public class ExpressionEvaluatingRequestHandlerAdvice extends AbstractRequestHan
}
}
private Object evaluateExpression(Message<?> message,
ExpressionEvaluatingMessageProcessor<Object> expressionEvaluatingMessageProcessor,
MessageChannel resultChannel, boolean propagateEvaluationFailure) throws Exception {
private void evaluateSuccessExpression(Message<?> message) throws Exception {
Object evalResult;
boolean evaluationFailed = false;
try {
evalResult = expressionEvaluatingMessageProcessor.processMessage(message);
evalResult = this.onSuccessExpression.getValue(this.prepareEvaluationContextToUse(null), message);
}
catch (Exception e) {
evalResult = e;
evaluationFailed = true;
}
if (evalResult != null && resultChannel != null) {
message = MessageBuilder.fromMessage(message)
.setHeader(MessageHeaders.POSTPROCESS_RESULT, evalResult)
.build();
this.messagingTemplate.send(resultChannel, message);
if (evalResult != null && this.successChannel != null) {
AdviceMessage resultMessage = new AdviceMessage(evalResult, message);
this.messagingTemplate.send(this.successChannel, resultMessage);
}
if (evaluationFailed && propagateEvaluationFailure) {
if (evaluationFailed && this.propagateOnSuccessEvaluationFailures) {
throw (Exception) evalResult;
}
}
private Object evaluateFailureExpression(Message<?> message, Exception exception) throws Exception {
Object evalResult;
try {
evalResult = this.onFailureExpression.getValue(this.prepareEvaluationContextToUse(exception), message);
}
catch (Exception e) {
evalResult = e;
logger.error("Failure expression evaluation failed for " + message + ": " + e.getMessage());
}
if (evalResult != null && this.failureChannel != null) {
MessagingException messagingException = new MessageHandlingExpressionEvaluatingAdviceException(message,
"Handler Failed", exception, evalResult);
ErrorMessage resultMessage = new ErrorMessage(messagingException);
this.messagingTemplate.send(this.failureChannel, resultMessage);
}
return evalResult;
}
protected StandardEvaluationContext createEvaluationContext(){
if (this.getBeanFactory() != null) {
return ExpressionUtils.createStandardEvaluationContext(new BeanFactoryResolver(this.getBeanFactory()),
this.getConversionService());
}
else {
return ExpressionUtils.createStandardEvaluationContext(this.getConversionService());
}
}
/**
* If we don't need variables (i.e., exception is null)
* we can use a singleton context; otherwise we need a new one each time.
* @param exception
* @return The context.
*/
private EvaluationContext prepareEvaluationContextToUse(Exception exception) {
EvaluationContext evaluationContextToUse;
if (exception != null) {
evaluationContextToUse = this.createEvaluationContext();
evaluationContextToUse.setVariable("exception", exception);
}
else {
if (this.evaluationContext == null) {
this.evaluationContext = this.createEvaluationContext();
}
evaluationContextToUse = this.evaluationContext;
}
return evaluationContextToUse;
}
public static class MessageHandlingExpressionEvaluatingAdviceException extends MessagingException {
private static final long serialVersionUID = 1L;
private final Object evaluationResult;
public MessageHandlingExpressionEvaluatingAdviceException(Message<?> message, String description,
Throwable cause, Object evaluationResult) {
super(message, description, cause);
this.evaluationResult = evaluationResult;
}
public Object getEvaluationResult() {
return evaluationResult;
}
}
}

View File

@@ -0,0 +1,53 @@
/*
* Copyright 2002-2012 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.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.message;
import java.util.Map;
import org.springframework.integration.Message;
/**
* A message implementation that is produced by an advice after
* successful message handling.
* Contains the result of the expression evaluation in the payload
* and the original message that the advice passed to the
* handler.
* .
* @author Gary Russell
* @since 2.2
*/
public class AdviceMessage extends GenericMessage<Object> {
private static final long serialVersionUID = 1L;
private final Message<?> inputMessage;
public AdviceMessage(Object payload, Message<?> inputMessage) {
super(payload);
this.inputMessage = inputMessage;
}
public AdviceMessage(Object payload, Map<String, Object> headers, Message<?> inputMessage) {
super(payload, headers);
this.inputMessage = inputMessage;
}
public Message<?> getInputMessage() {
return inputMessage;
}
}

View File

@@ -28,13 +28,14 @@ import java.util.concurrent.atomic.AtomicInteger;
import org.aopalliance.aop.Advice;
import org.junit.Test;
import org.springframework.expression.spel.standard.SpelExpressionParser;
import org.springframework.integration.Message;
import org.springframework.integration.MessageHandlingException;
import org.springframework.integration.MessageHeaders;
import org.springframework.integration.MessagingException;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.core.PollableChannel;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
import org.springframework.integration.handler.advice.ExpressionEvaluatingRequestHandlerAdvice.MessageHandlingExpressionEvaluatingAdviceException;
import org.springframework.integration.message.AdviceMessage;
import org.springframework.integration.message.GenericMessage;
import org.springframework.retry.RecoveryCallback;
import org.springframework.retry.RetryContext;
@@ -72,9 +73,11 @@ public class AdvisedMessageHandlerTests {
PollableChannel successChannel = new QueueChannel();
PollableChannel failureChannel = new QueueChannel();
ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice(
"'foo'", successChannel,
"'bar'", failureChannel);
ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice();
advice.setSuccessChannel(successChannel);
advice.setFailureChannel(failureChannel);
advice.setOnSuccessExpression("'foo'");
advice.setOnFailureExpression("'bar:' + #exception.cause.message");
List<Advice> adviceChain = new ArrayList<Advice>();
adviceChain.add(advice);
@@ -89,8 +92,8 @@ public class AdvisedMessageHandlerTests {
Message<?> success = successChannel.receive(1000);
assertNotNull(success);
assertEquals("Hello, world!", success.getPayload());
assertEquals("foo", success.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT));
assertEquals("Hello, world!", ((AdviceMessage) success).getInputMessage().getPayload());
assertEquals("foo", success.getPayload());
// advice with failure, not trapped
doFail.set(true);
@@ -104,16 +107,16 @@ public class AdvisedMessageHandlerTests {
Message<?> failure = failureChannel.receive(1000);
assertNotNull(failure);
assertEquals("Hello, world!", failure.getPayload());
assertEquals("bar", failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT));
assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload());
assertEquals("bar:qux", ((MessageHandlingExpressionEvaluatingAdviceException) failure.getPayload()).getEvaluationResult());
// advice with failure, trapped
advice.setTrapException(true);
handler.handleMessage(message);
failure = failureChannel.receive(1000);
assertNotNull(failure);
assertEquals("Hello, world!", failure.getPayload());
assertEquals("bar", failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT));
assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload());
assertEquals("bar:qux", ((MessageHandlingExpressionEvaluatingAdviceException) failure.getPayload()).getEvaluationResult());
assertNull(replies.receive(1));
// advice with failure, eval is result
@@ -121,12 +124,12 @@ public class AdvisedMessageHandlerTests {
handler.handleMessage(message);
failure = failureChannel.receive(1000);
assertNotNull(failure);
assertEquals("Hello, world!", failure.getPayload());
assertEquals("bar", failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT));
assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload());
assertEquals("bar:qux", ((MessageHandlingExpressionEvaluatingAdviceException) failure.getPayload()).getEvaluationResult());
reply = replies.receive(1000);
assertNotNull(reply);
assertEquals("bar", reply.getPayload());
assertEquals("bar:qux", reply.getPayload());
}
@@ -148,9 +151,11 @@ public class AdvisedMessageHandlerTests {
PollableChannel successChannel = new QueueChannel();
PollableChannel failureChannel = new QueueChannel();
ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice(
new SpelExpressionParser().parseExpression("1/0"), successChannel,
new SpelExpressionParser().parseExpression("1/0"), failureChannel);
ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice();
advice.setSuccessChannel(successChannel);
advice.setFailureChannel(failureChannel);
advice.setOnSuccessExpression("1/0");
advice.setOnFailureExpression("1/0");
List<Advice> adviceChain = new ArrayList<Advice>();
adviceChain.add(advice);
@@ -165,9 +170,9 @@ public class AdvisedMessageHandlerTests {
Message<?> success = successChannel.receive(1000);
assertNotNull(success);
assertEquals("Hello, world!", success.getPayload());
assertEquals(MessageHandlingException.class, success.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT).getClass());
assertEquals("Expression evaluation failed: 1/0", ((Exception) success.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT)).getMessage());
assertEquals("Hello, world!", ((AdviceMessage) success).getInputMessage().getPayload());
assertEquals(ArithmeticException.class, success.getPayload().getClass());
assertEquals("/ by zero", ((Exception) success.getPayload()).getMessage());
// propagate failing advice with success
advice.setPropagateEvaluationFailures(true);
@@ -176,16 +181,16 @@ public class AdvisedMessageHandlerTests {
fail("Expected Exception");
}
catch (MessageHandlingException e) {
assertEquals("Expression evaluation failed: 1/0", e.getMessage());
assertEquals("/ by zero", e.getCause().getMessage());
}
reply = replies.receive(1);
assertNull(reply);
success = successChannel.receive(1000);
assertNotNull(success);
assertEquals("Hello, world!", success.getPayload());
assertEquals(MessageHandlingException.class, success.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT).getClass());
assertEquals("Expression evaluation failed: 1/0", ((Exception) success.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT)).getMessage());
assertEquals("Hello, world!", ((AdviceMessage) success).getInputMessage().getPayload());
assertEquals(ArithmeticException.class, success.getPayload().getClass());
assertEquals("/ by zero", ((Exception) success.getPayload()).getMessage());
}
@@ -207,9 +212,11 @@ public class AdvisedMessageHandlerTests {
PollableChannel successChannel = new QueueChannel();
PollableChannel failureChannel = new QueueChannel();
ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice(
new SpelExpressionParser().parseExpression("1/0"), successChannel,
new SpelExpressionParser().parseExpression("1/0"), failureChannel);
ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice();
advice.setSuccessChannel(successChannel);
advice.setFailureChannel(failureChannel);
advice.setOnSuccessExpression("1/0");
advice.setOnFailureExpression("1/0");
List<Advice> adviceChain = new ArrayList<Advice>();
adviceChain.add(advice);
@@ -229,9 +236,9 @@ public class AdvisedMessageHandlerTests {
Message<?> failure = failureChannel.receive(1000);
assertNotNull(failure);
assertEquals("Hello, world!", failure.getPayload());
assertEquals(MessageHandlingException.class, failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT).getClass());
assertEquals("Expression evaluation failed: 1/0", ((Exception) failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT)).getMessage());
assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload());
assertEquals(MessageHandlingExpressionEvaluatingAdviceException.class, failure.getPayload().getClass());
assertEquals("qux", ((Exception) failure.getPayload()).getCause().getCause().getMessage());
// propagate failing advice with failure; expect original exception
advice.setPropagateEvaluationFailures(true);
@@ -247,9 +254,9 @@ public class AdvisedMessageHandlerTests {
failure = failureChannel.receive(1000);
assertNotNull(failure);
assertEquals("Hello, world!", failure.getPayload());
assertEquals(MessageHandlingException.class, failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT).getClass());
assertEquals("Expression evaluation failed: 1/0", ((Exception) failure.getHeaders().get(MessageHeaders.POSTPROCESS_RESULT)).getMessage());
assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload());
assertEquals(MessageHandlingExpressionEvaluatingAdviceException.class, failure.getPayload().getClass());
assertEquals("qux", ((Exception) failure.getPayload()).getCause().getCause().getMessage());
}

View File

@@ -369,46 +369,24 @@ Caused by: java.lang.RuntimeException: foo
<classname>org.springframework.integration.handler.advice.ExpressionEvaluatingRequestHandlerAdvice</classname>.
This advice is more general than the other two advices. It provides a mechanism to evaluate an expression on the
original inbound message sent to the endpoint. Separate expressions are available to be evaluated, either after
success, or failure. Optionally, the original message, together with the result of the evaluation in a
header, can be sent to a message channel.
success, or failure. Optionally, a message containing the evaluation result, together with the input message,
can be sent to a message channel.
</para>
<para>
A typical use case for this advice might be with an &lt;ftp:outbound-channel-adapter /&gt;, perhaps to move
the file to one directory if the transfer was successful, or to another directory if it fails:
</para>
<para><programlisting><![CDATA[<int-ftp:outbound-channel-adapter id="ftpOutbound" cache-sessions="false"
channel="ftpChannel"
remote-directory="."
session-factory="ftpClientFactory">
<int-ftp:request-handler-advice-chain>
<bean class="org.springframework.integration.handler.advice.ExpressionEvaluatingRequestHandlerAdvice">
<constructor-arg name="onSuccessExpression" value="payload.renameTo('/tmp/good/' + payload.name)"/>
<constructor-arg name="successChannel" ref="successChannel" />
<constructor-arg name="onFailureExpression" value="payload.renameTo('/tmp/bad/' + payload.name)"/>
<constructor-arg name="failureChannel" ref="failureChannel" />
</bean>
</int-ftp:request-handler-advice-chain>
</int-ftp:outbound-channel-adapter>
DEBUG [main] preSend on channel 'ftpChannel', message: [Payload=target/toSend/b.txt][Headers=...]
DEBUG [main] Connected to server [localhost:21]
INFO [main] File has been successfully transfered to: ./b.txt.writing
INFO [main] File has been successfully renamed from: ./b.txt.writing to ./b.txt
DEBUG [main] preSend on channel 'successChannel', message: [Payload=target/toSend/b.txt][Headers={..., postProcessResult=true}]
...
DEBUG [main] preSend on channel 'ftpChannel', message: [Payload=target/toSend/a.txt][Headers=...]
DEBUG [main] Connected to server [localhost:21]
...
DEBUG [main] preSend on channel 'failureChannel', message: [Payload=target/toSend/a.txt][Headers={..., postProcessResult=true}]]]></programlisting></para>
<para>
As you can see, in the first case, the successful transfer resulted in a message sent to
<emphasis>successChannel</emphasis>; the second case shows the result being sent to
<emphasis>failureChannel</emphasis>. In both cases, you can see the <emphasis>postProcessResult</emphasis>
header containing the result of the evaluation (true because the rename operations were successful -
<classname>java.io.File.renameTo(...)</classname> returns a <emphasis>boolean</emphasis>).
The Advice has properties to set an expression when successful, an expression for failures,
and corresponding channels for each. For the successful case, the message sent to the
<emphasis>successChannel</emphasis> is an <classname>AdviceMessage</classname>, with
the payload being the result of the expression evaluation, and an additional property
<code>inputMessage</code> which contains the original message sent to the handler. A message
sent to the <emphasis>failureChannel</emphasis> (when the handler throws an excecption)
is an ErrorMessage with a payload of <classname>MessageHandlingExpressionEvaluatingAdviceException</classname>.
Like all <classname>MessagingException</classname>s, this payload has <code>failedMessage</code>
and <code>cause</code> properties, as well as an additional property <code>evaluationResult</code>,
containing the result of the expression evaluation.
</para>
</section>
</section>