INT-3998: Channel Late-Binding for WireTap
JIRA: https://jira.spring.io/browse/INT-3998 Do not modify interceptors from the `AbstractMessageChannel`. They must be declared as beans, too. Fix [UnusedImport] in the `EnableIntegrationTests` Doc Polishing
This commit is contained in:
committed by
Gary Russell
parent
d8e4a9d114
commit
882f4c017e
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2015 the original author or authors.
|
||||
* Copyright 2002-2016 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.
|
||||
@@ -19,9 +19,13 @@ package org.springframework.integration.channel.interceptor;
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.beans.BeansException;
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.beans.factory.BeanFactoryAware;
|
||||
import org.springframework.context.Lifecycle;
|
||||
import org.springframework.integration.channel.ChannelInterceptorAware;
|
||||
import org.springframework.integration.core.MessageSelector;
|
||||
import org.springframework.integration.support.channel.BeanFactoryChannelResolver;
|
||||
import org.springframework.jmx.export.annotation.ManagedAttribute;
|
||||
import org.springframework.jmx.export.annotation.ManagedOperation;
|
||||
import org.springframework.jmx.export.annotation.ManagedResource;
|
||||
@@ -37,13 +41,17 @@ import org.springframework.util.Assert;
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@ManagedResource
|
||||
public class WireTap extends ChannelInterceptorAdapter implements Lifecycle, VetoCapableInterceptor {
|
||||
public class WireTap extends ChannelInterceptorAdapter
|
||||
implements Lifecycle, VetoCapableInterceptor, BeanFactoryAware {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(WireTap.class);
|
||||
|
||||
private final MessageChannel channel;
|
||||
private volatile MessageChannel channel;
|
||||
|
||||
private volatile String channelName;
|
||||
|
||||
private volatile long timeout = 0;
|
||||
|
||||
@@ -51,10 +59,11 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle, Vet
|
||||
|
||||
private volatile boolean running = true;
|
||||
|
||||
private BeanFactory beanFactory;
|
||||
|
||||
|
||||
/**
|
||||
* Create a new wire tap with <em>no</em> {@link MessageSelector}.
|
||||
*
|
||||
* @param channel the MessageChannel to which intercepted messages will be sent
|
||||
*/
|
||||
public WireTap(MessageChannel channel) {
|
||||
@@ -63,7 +72,6 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle, Vet
|
||||
|
||||
/**
|
||||
* Create a new wire tap with the provided {@link MessageSelector}.
|
||||
*
|
||||
* @param channel the channel to which intercepted messages will be sent
|
||||
* @param selector the selector that must accept a message for it to be
|
||||
* sent to the intercepting channel
|
||||
@@ -74,16 +82,46 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle, Vet
|
||||
this.selector = selector;
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a new wire tap based on the MessageChannel name and
|
||||
* with <em>no</em> {@link MessageSelector}.
|
||||
* @param channelName the name of the target MessageChannel
|
||||
* to which intercepted messages will be sent
|
||||
* @since 4.3
|
||||
*/
|
||||
public WireTap(String channelName) {
|
||||
this(channelName, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a new wire tap with the provided {@link MessageSelector}.
|
||||
* @param channelName the name of the target MessageChannel
|
||||
* to which intercepted messages will be sent.
|
||||
* @param selector the selector that must accept a message for it to be
|
||||
* sent to the intercepting channel
|
||||
* @since 4.3
|
||||
*/
|
||||
public WireTap(String channelName, MessageSelector selector) {
|
||||
Assert.hasText(channelName, "channelName must not be empty");
|
||||
this.channelName = channelName;
|
||||
this.selector = selector;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify the timeout value for sending to the intercepting target.
|
||||
*
|
||||
* @param timeout the timeout in milliseconds
|
||||
*/
|
||||
public void setTimeout(long timeout) {
|
||||
this.timeout = timeout;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
|
||||
if (this.beanFactory == null) {
|
||||
this.beanFactory = beanFactory;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Check whether the wire tap is currently running.
|
||||
*/
|
||||
@@ -118,18 +156,19 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle, Vet
|
||||
*/
|
||||
@Override
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (this.channel.equals(channel)) {
|
||||
MessageChannel wireTapChannel = getChannel();
|
||||
if (wireTapChannel.equals(channel)) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("WireTap is refusing to intercept its own channel '" + this.channel + "'");
|
||||
logger.debug("WireTap is refusing to intercept its own channel '" + wireTapChannel + "'");
|
||||
}
|
||||
return message;
|
||||
}
|
||||
if (this.running && (this.selector == null || this.selector.accept(message))) {
|
||||
boolean sent = (this.timeout >= 0)
|
||||
? this.channel.send(message, this.timeout)
|
||||
: this.channel.send(message);
|
||||
? wireTapChannel.send(message, this.timeout)
|
||||
: wireTapChannel.send(message);
|
||||
if (!sent && logger.isWarnEnabled()) {
|
||||
logger.warn("failed to send message to WireTap channel '" + this.channel + "'");
|
||||
logger.warn("failed to send message to WireTap channel '" + wireTapChannel + "'");
|
||||
}
|
||||
}
|
||||
return message;
|
||||
@@ -137,7 +176,20 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle, Vet
|
||||
|
||||
@Override
|
||||
public boolean shouldIntercept(String beanName, ChannelInterceptorAware channel) {
|
||||
return !this.channel.equals(channel);
|
||||
return !getChannel().equals(channel);
|
||||
}
|
||||
|
||||
private MessageChannel getChannel() {
|
||||
if (this.channelName != null) {
|
||||
synchronized (this) {
|
||||
if (this.channelName != null) {
|
||||
this.channel = new BeanFactoryChannelResolver(this.beanFactory)
|
||||
.resolveDestination(this.channelName);
|
||||
this.channelName = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
return this.channel;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -25,11 +25,13 @@ import org.junit.Test;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.core.MessageSelector;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
public class WireTapTests {
|
||||
|
||||
@@ -73,7 +75,7 @@ public class WireTapTests {
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void wireTapTargetMustNotBeNull() {
|
||||
new WireTap(null);
|
||||
new WireTap((MessageChannel) null);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -98,7 +100,7 @@ public class WireTapTests {
|
||||
mainChannel.addInterceptor(new WireTap(secondaryChannel));
|
||||
String headerName = "testAttribute";
|
||||
Message<String> message = MessageBuilder.withPayload("testing")
|
||||
.setHeader(headerName, new Integer(123)).build();
|
||||
.setHeader(headerName, 123).build();
|
||||
mainChannel.send(message);
|
||||
Message<?> original = mainChannel.receive(0);
|
||||
Message<?> intercepted = secondaryChannel.receive(0);
|
||||
|
||||
@@ -185,6 +185,9 @@ public class EnableIntegrationTests {
|
||||
@Autowired
|
||||
private QueueChannel output;
|
||||
|
||||
@Autowired
|
||||
private QueueChannel wireTapFromOutput;
|
||||
|
||||
@Autowired
|
||||
private PollableChannel publishedChannel;
|
||||
|
||||
@@ -299,6 +302,7 @@ public class EnableIntegrationTests {
|
||||
invocation.callRealMethod();
|
||||
return null;
|
||||
}
|
||||
|
||||
}).when(logger).debug("Received no Message during the poll, returning 'false'");
|
||||
new DirectFieldAccessor(this.serviceActivatorEndpoint).setPropertyValue("logger", logger);
|
||||
|
||||
@@ -345,6 +349,10 @@ public class EnableIntegrationTests {
|
||||
assertNotNull(receive);
|
||||
assertEquals("FOO", receive.getPayload());
|
||||
|
||||
receive = this.wireTapFromOutput.receive(10000);
|
||||
assertNotNull(receive);
|
||||
assertEquals("FOO", receive.getPayload());
|
||||
|
||||
MessageHistory messageHistory = receive.getHeaders().get(MessageHistory.HEADER_NAME, MessageHistory.class);
|
||||
assertNotNull(messageHistory);
|
||||
String messageHistoryString = messageHistory.toString();
|
||||
@@ -722,8 +730,20 @@ public class EnableIntegrationTests {
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public WireTap wireTapFromOutputInterceptor() {
|
||||
return new WireTap("wireTapFromOutput");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public PollableChannel output() {
|
||||
QueueChannel queueChannel = new QueueChannel();
|
||||
queueChannel.addInterceptor(wireTapFromOutputInterceptor());
|
||||
return queueChannel;
|
||||
}
|
||||
|
||||
@Bean
|
||||
public PollableChannel wireTapFromOutput() {
|
||||
return new QueueChannel();
|
||||
}
|
||||
|
||||
|
||||
@@ -788,6 +788,30 @@ That way, the framework will ask the interceptor if it's OK to intercept each ch
|
||||
You can also add runtime protection in the interceptor methods that ensures that the channel is not one that is referenced by the interceptor.
|
||||
The `WireTap` uses both of these techniques.
|
||||
|
||||
Starting with _version 4.3_, the `WireTap` has additional constructors that take a `channelName` instead of a
|
||||
`MessageChannel` instance.
|
||||
This can be convenient for Java Configuration and when channel auto-creation logic is being used.
|
||||
The target `MessageChannel` bean is resolved from the provided `channelName` later, on the first interaction with the
|
||||
interceptor.
|
||||
|
||||
IMPORTANT: Channel resolution requires a `BeanFactory` so the wire tap instance must be a Spring-managed bean.
|
||||
|
||||
This _late-binding_ approach also allows simplification of typical wire-tapping patterns with Java DSL configuration:
|
||||
[source,java]
|
||||
----
|
||||
@Bean
|
||||
public PollableChannel myChannel() {
|
||||
return MessageChannels.queue()
|
||||
.wireTap("loggingFlow.input")
|
||||
.get();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public IntegrationFlow loggingFlow() {
|
||||
return f -> f.log();
|
||||
}
|
||||
----
|
||||
|
||||
[[conditional-wiretap]]
|
||||
===== Conditional Wire Taps
|
||||
|
||||
|
||||
@@ -206,3 +206,9 @@ See <<annotations>> for more information.
|
||||
|
||||
The XMPP Extensions (XEP) are now supported by the XMPP channel adapters.
|
||||
See <<xmpp-extensions>> for more information.
|
||||
|
||||
==== WireTap Late Binding
|
||||
|
||||
The `WireTap` `ChannelInterceptor` now can accept a `channelName` which is resolved to the target `MessageChannel`
|
||||
later, during the first active interceptor operation.
|
||||
See <<channel-wiretap>> for more information.
|
||||
|
||||
Reference in New Issue
Block a user