Message payload can no longer be set. MessageTransformer's transform() method now returns a Message (instead of void). ChannelInterceptor preSend() and postReceive() methods now return a Message instead of boolean.
This commit is contained in:
@@ -54,7 +54,7 @@ public class DefaultErrorChannel extends RendezvousChannel {
|
||||
* Even if the error channel has no subscribers, errors are at least visible at debug level.
|
||||
*/
|
||||
@Override
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
String errorMessage = "Error Received. Message: " + message.toString();
|
||||
Object payload = message.getPayload();
|
||||
@@ -65,7 +65,7 @@ public class DefaultErrorChannel extends RendezvousChannel {
|
||||
logger.debug(errorMessage);
|
||||
}
|
||||
}
|
||||
return true;
|
||||
return message;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -110,7 +110,8 @@ public abstract class AbstractMessageChannel implements MessageChannel, BeanName
|
||||
* time or the sending thread is interrupted.
|
||||
*/
|
||||
public final boolean send(Message<?> message, long timeout) {
|
||||
if (!this.interceptors.preSend(message, this)) {
|
||||
message = this.interceptors.preSend(message, this);
|
||||
if (message == null) {
|
||||
return false;
|
||||
}
|
||||
boolean sent = this.doSend(message, timeout);
|
||||
@@ -147,7 +148,7 @@ public abstract class AbstractMessageChannel implements MessageChannel, BeanName
|
||||
return null;
|
||||
}
|
||||
Message<?> message = this.doReceive(timeout);
|
||||
this.interceptors.postReceive(message, this);
|
||||
message = this.interceptors.postReceive(message, this);
|
||||
return message;
|
||||
}
|
||||
|
||||
@@ -194,16 +195,17 @@ public abstract class AbstractMessageChannel implements MessageChannel, BeanName
|
||||
return this.interceptors.add(interceptor);
|
||||
}
|
||||
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("preSend on channel '" + channel + "', message: " + message);
|
||||
}
|
||||
for (ChannelInterceptor interceptor : interceptors) {
|
||||
if (!interceptor.preSend(message, channel)) {
|
||||
return false;
|
||||
message = interceptor.preSend(message, channel);
|
||||
if (message == null) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
return message;
|
||||
}
|
||||
|
||||
public void postSend(Message<?> message, MessageChannel channel, boolean sent) {
|
||||
@@ -227,7 +229,7 @@ public abstract class AbstractMessageChannel implements MessageChannel, BeanName
|
||||
return true;
|
||||
}
|
||||
|
||||
public void postReceive(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
if (message != null && logger.isDebugEnabled()) {
|
||||
logger.debug("postReceive on channel '" + channel + "', message: " + message);
|
||||
}
|
||||
@@ -235,8 +237,12 @@ public abstract class AbstractMessageChannel implements MessageChannel, BeanName
|
||||
logger.trace("postReceive on channel '" + channel + "', message is null");
|
||||
}
|
||||
for (ChannelInterceptor interceptor : interceptors) {
|
||||
interceptor.postReceive(message, channel);
|
||||
message = interceptor.postReceive(message, channel);
|
||||
if (message == null) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
return message;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -27,12 +27,12 @@ import org.springframework.integration.message.Message;
|
||||
*/
|
||||
public interface ChannelInterceptor {
|
||||
|
||||
boolean preSend(Message<?> message, MessageChannel channel);
|
||||
Message<?> preSend(Message<?> message, MessageChannel channel);
|
||||
|
||||
void postSend(Message<?> message, MessageChannel channel, boolean sent);
|
||||
|
||||
boolean preReceive(MessageChannel channel);
|
||||
|
||||
void postReceive(Message<?> message, MessageChannel channel);
|
||||
Message<?> postReceive(Message<?> message, MessageChannel channel);
|
||||
|
||||
}
|
||||
|
||||
@@ -28,8 +28,8 @@ import org.springframework.integration.message.Message;
|
||||
*/
|
||||
public class ChannelInterceptorAdapter implements ChannelInterceptor {
|
||||
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
return true;
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
return message;
|
||||
}
|
||||
|
||||
public void postSend(Message<?> message, MessageChannel channel, boolean sent) {
|
||||
@@ -39,7 +39,8 @@ public class ChannelInterceptorAdapter implements ChannelInterceptor {
|
||||
return true;
|
||||
}
|
||||
|
||||
public void postReceive(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -42,14 +42,14 @@ public class MessageSelectingInterceptor extends ChannelInterceptorAdapter {
|
||||
|
||||
|
||||
@Override
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
for (MessageSelector selector : this.selectors) {
|
||||
if (!selector.accept(message)) {
|
||||
throw new MessageDeliveryException(message,
|
||||
"selector '" + selector + "' did not accept message");
|
||||
}
|
||||
}
|
||||
return true;
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -54,11 +54,11 @@ public class MessageStoringInterceptor extends ChannelInterceptorAdapter {
|
||||
|
||||
|
||||
@Override
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (message != null) {
|
||||
this.messageStore.put(message.getId(), message);
|
||||
}
|
||||
return true;
|
||||
return message;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -84,10 +84,11 @@ public class MessageStoringInterceptor extends ChannelInterceptorAdapter {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void postReceive(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
if (message != null) {
|
||||
this.messageStore.remove(message.getId());
|
||||
}
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -98,7 +98,7 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (this.running && this.selectorsAccept(message)) {
|
||||
Message<?> duplicate = new GenericMessage<Object>(message.getPayload(), message.getHeader());
|
||||
duplicate.getHeader().setAttribute(ORIGINAL_MESSAGE_ID_KEY, message.getId());
|
||||
@@ -109,7 +109,7 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle {
|
||||
}
|
||||
}
|
||||
}
|
||||
return true;
|
||||
return message;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -31,8 +31,6 @@ public interface Message<T> extends Serializable {
|
||||
|
||||
T getPayload();
|
||||
|
||||
void setPayload(T payload);
|
||||
|
||||
boolean isExpired();
|
||||
|
||||
void copyHeader(MessageHeader header, boolean overwriteExistingValues);
|
||||
|
||||
@@ -22,6 +22,7 @@ import java.util.Properties;
|
||||
import org.springframework.integration.handler.MessageHandler;
|
||||
import org.springframework.integration.handler.annotation.AnnotationMethodMessageMapper;
|
||||
import org.springframework.integration.message.DefaultMessageMapper;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.message.MessageMapper;
|
||||
import org.springframework.integration.message.MessagingException;
|
||||
@@ -30,8 +31,7 @@ import org.springframework.integration.util.AbstractMethodInvokingAdapter;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class AnnotationMethodTransformerAdapter extends AbstractMethodInvokingAdapter
|
||||
implements MessageTransformer, MessageHandler {
|
||||
public class AnnotationMethodTransformerAdapter extends AbstractMethodInvokingAdapter implements MessageHandler {
|
||||
|
||||
private volatile MessageMapper mapper;
|
||||
|
||||
@@ -49,12 +49,12 @@ public class AnnotationMethodTransformerAdapter extends AbstractMethodInvokingAd
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public void transform(Message message) {
|
||||
public Message<?> transform(Message<?> message) {
|
||||
if (!this.isInitialized()) {
|
||||
this.afterPropertiesSet();
|
||||
}
|
||||
if (message.getPayload() == null) {
|
||||
return;
|
||||
return message;
|
||||
}
|
||||
Object param = (this.methodExpectsMessage) ? message : this.mapper.mapMessage(message);
|
||||
try {
|
||||
@@ -78,7 +78,7 @@ public class AnnotationMethodTransformerAdapter extends AbstractMethodInvokingAd
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("MessageTransformer returned a null result");
|
||||
}
|
||||
return;
|
||||
return null;
|
||||
}
|
||||
if (result instanceof Properties && !(message.getPayload() instanceof Properties)) {
|
||||
Properties propertiesToSet = (Properties) result;
|
||||
@@ -94,17 +94,17 @@ public class AnnotationMethodTransformerAdapter extends AbstractMethodInvokingAd
|
||||
}
|
||||
}
|
||||
else {
|
||||
message.setPayload(result);
|
||||
return new GenericMessage(result, message.getHeader());
|
||||
}
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new MessagingException(message, "failed to transform message payload", e);
|
||||
}
|
||||
}
|
||||
|
||||
public Message<?> handle(Message<?> message) {
|
||||
this.transform(message);
|
||||
return message;
|
||||
}
|
||||
|
||||
public Message<?> handle(Message<?> message) {
|
||||
return this.transform(message);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -25,6 +25,6 @@ import org.springframework.integration.message.Message;
|
||||
*/
|
||||
public interface MessageTransformer {
|
||||
|
||||
void transform(Message<?> message);
|
||||
Message<?> transform(Message<?> message);
|
||||
|
||||
}
|
||||
|
||||
@@ -48,10 +48,14 @@ public class MessageTransformerChain implements MessageTransformer {
|
||||
this.transformers.addAll(transformers);
|
||||
}
|
||||
|
||||
public void transform(Message<?> message) {
|
||||
for (MessageTransformer next : transformers) {
|
||||
next.transform(message);
|
||||
public Message<?> transform(Message<?> message) {
|
||||
for (MessageTransformer next : this.transformers) {
|
||||
message = next.transform(message);
|
||||
if (message == null) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -48,18 +48,19 @@ public class MessageTransformingChannelInterceptor extends ChannelInterceptorAda
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean preSend(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (this.transformOnSend) {
|
||||
this.transfomer.transform(message);
|
||||
message = this.transfomer.transform(message);
|
||||
}
|
||||
return true;
|
||||
return message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void postReceive(Message<?> message, MessageChannel channel) {
|
||||
public Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
if (!this.transformOnSend) {
|
||||
this.transfomer.transform(message);
|
||||
message = this.transfomer.transform(message);
|
||||
}
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
|
||||
package org.springframework.integration.transformer;
|
||||
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.message.MessagingException;
|
||||
import org.springframework.integration.util.AbstractMethodInvokingAdapter;
|
||||
@@ -26,9 +27,10 @@ import org.springframework.integration.util.AbstractMethodInvokingAdapter;
|
||||
public class PayloadTransformerAdapter extends AbstractMethodInvokingAdapter implements MessageTransformer {
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public void transform(Message message) {
|
||||
public Message<?> transform(Message<?> message) {
|
||||
try {
|
||||
message.setPayload(this.invokeMethod(message.getPayload()));
|
||||
Object result = this.invokeMethod(message.getPayload());
|
||||
return new GenericMessage(result, message.getHeader());
|
||||
} catch (Exception e) {
|
||||
throw new MessagingException(message, "failed to transform message payload", e);
|
||||
}
|
||||
|
||||
@@ -33,8 +33,7 @@ public class TransformerMessageHandlerAdapter implements MessageHandler {
|
||||
|
||||
|
||||
public Message<?> handle(Message<?> message) {
|
||||
this.transformer.transform(message);
|
||||
return message;
|
||||
return this.transformer.transform(message);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user