diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/WireTap.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/WireTap.java index 6f511a3c02..e5a61ccf11 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/WireTap.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/WireTap.java @@ -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 no {@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 no {@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; } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/WireTapTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/WireTapTests.java index 07a02a595a..ac388d0be8 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/WireTapTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/WireTapTests.java @@ -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 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); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java index e9d229879e..220fa06812 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java @@ -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(); } diff --git a/src/reference/asciidoc/channel.adoc b/src/reference/asciidoc/channel.adoc index 47bc7ff036..a025eb3066 100644 --- a/src/reference/asciidoc/channel.adoc +++ b/src/reference/asciidoc/channel.adoc @@ -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 diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index d9e3622cf2..a0c786b251 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -206,3 +206,9 @@ See <> for more information. The XMPP Extensions (XEP) are now supported by the XMPP channel adapters. See <> 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 <> for more information.