diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java index 625da7f108..78bf9e672a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java @@ -20,6 +20,7 @@ import java.util.Collection; import java.util.HashSet; import java.util.LinkedHashSet; import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.concurrent.Executor; import java.util.function.Consumer; @@ -42,6 +43,7 @@ import org.springframework.integration.channel.ReactiveChannel; import org.springframework.integration.channel.interceptor.WireTap; import org.springframework.integration.config.ConsumerEndpointFactoryBean; import org.springframework.integration.config.SourcePollingChannelAdapterFactoryBean; +import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.core.GenericSelector; import org.springframework.integration.core.MessageSelector; import org.springframework.integration.dsl.channel.MessageChannelSpec; @@ -412,7 +414,7 @@ public abstract class IntegrationFlowDefinition The full request {@link Message} will be logged. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public B log() { return log(LoggingHandler.Level.INFO); @@ -2374,6 +2377,7 @@ public abstract class IntegrationFlowDefinition The full request {@link Message} will be logged. * @param level the {@link LoggingHandler.Level}. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public B log(LoggingHandler.Level level) { return log(level, (String) null); @@ -2386,6 +2390,7 @@ public abstract class IntegrationFlowDefinition The full request {@link Message} will be logged. * @param category the logging category to use. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public B log(String category) { return log(LoggingHandler.Level.INFO, category); @@ -2399,6 +2404,7 @@ public abstract class IntegrationFlowDefinition the expected payload type. * against the request {@link Message}. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public

B log(Function, Object> function) { Assert.notNull(function); @@ -2444,6 +2452,7 @@ public abstract class IntegrationFlowDefinition the expected payload type. * against the request {@link Message}. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public

B log(LoggingHandler.Level level, Function, Object> function) { return log(level, null, function); @@ -2507,6 +2519,7 @@ public abstract class IntegrationFlowDefinition the expected payload type. * against the request {@link Message}. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public

B log(String category, Function, Object> function) { return log(LoggingHandler.Level.INFO, category, function); @@ -2523,6 +2536,7 @@ public abstract class IntegrationFlowDefinition the expected payload type. * against the request {@link Message}. * @return the current {@link IntegrationFlowDefinition}. + * @see #wireTap(WireTapSpec) */ public

B log(LoggingHandler.Level level, String category, Function, Object> function) { Assert.notNull(function); @@ -2540,6 +2554,7 @@ public abstract class IntegrationFlowDefinition lastComponent = this.integrationComponents.stream().reduce((first, second) -> second); + if (lastComponent.get() instanceof WireTapSpec) { +// channel(IntegrationContextUtils.NULL_CHANNEL_BEAN_NAME); + } + this.integrationFlow = new StandardIntegrationFlow(this.integrationComponents); } return this.integrationFlow; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java index 7d13fda9ba..cbb297c21b 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 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. @@ -16,7 +16,6 @@ package org.springframework.integration.dsl.channel; -import java.util.Arrays; import java.util.Collection; import java.util.Collections; @@ -104,10 +103,10 @@ public class WireTapSpec extends IntegrationComponentSpec @Override public Collection getComponentsToRegister() { if (this.selector != null) { - return Arrays.asList(this.selector, this.target); + return Collections.singleton(this.selector); } else { - return Collections.singletonList(this.target); + return null; } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java index 571663b7df..3e5fb96776 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 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. @@ -49,6 +49,7 @@ import org.springframework.integration.dsl.IntegrationFlowAdapter; import org.springframework.integration.dsl.IntegrationFlowDefinition; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.dsl.channel.MessageChannels; +import org.springframework.integration.handler.LoggingHandler; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -57,8 +58,7 @@ import org.springframework.messaging.handler.annotation.Header; import org.springframework.scheduling.TriggerContext; import org.springframework.stereotype.Component; import org.springframework.test.annotation.DirtiesContext; -import org.springframework.test.context.ContextConfiguration; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.junit4.SpringRunner; import org.springframework.util.StringUtils; /** @@ -66,8 +66,7 @@ import org.springframework.util.StringUtils; * * @since 5.0 */ -@ContextConfiguration -@RunWith(SpringJUnit4ClassRunner.class) +@RunWith(SpringRunner.class) @DirtiesContext public class FlowServiceTests { @@ -82,13 +81,12 @@ public class FlowServiceTests { private PollableChannel myFlowAdapterOutput; @Test - public void testFlowService() { + public void testFlowServiceAndLogAsLastNoError() { assertNotNull(this.myFlow); - QueueChannel replyChannel = new QueueChannel(); - this.input.send(MessageBuilder.withPayload("foo").setReplyChannel(replyChannel).build()); - Message receive = replyChannel.receive(1000); - assertNotNull(receive); - assertEquals("FOO", receive.getPayload()); + this.input.send(MessageBuilder.withPayload("foo").build()); + Object result = this.myFlow.resultOverLoggingHandler.get(); + assertNotNull(result); + assertEquals("FOO", result); } @Test @@ -119,7 +117,9 @@ public class FlowServiceTests { @Bean public IntegrationFlow testGateway() { - return f -> f.gateway("processChannel", g -> g.replyChannel("replyChannel")); + return f -> f.gateway("processChannel", g -> g.replyChannel("replyChannel")) + .log() + .bridge(null); } @Bean @@ -136,9 +136,15 @@ public class FlowServiceTests { @Component public static class MyFlow implements IntegrationFlow { + private final AtomicReference resultOverLoggingHandler = new AtomicReference<>(); + @Override public void configure(IntegrationFlowDefinition f) { - f.transform(String::toUpperCase); + f.transform(String::toUpperCase) + .log(LoggingHandler.Level.ERROR, m -> { + resultOverLoggingHandler.set(m.getPayload()); + return m; + }); } }