From 3be95190be279567528b3d1f43c4885a1ccfa6c7 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 11 Jan 2023 18:53:20 +0100 Subject: [PATCH] GH-2563 Add warning when SB is sending to input binding Resolves #2563 --- .../ImplicitFunctionBindingTests.java | 2 +- .../stream/function/StreamBridgeTests.java | 37 +++++++++++++------ .../cloud/stream/function/StreamBridge.java | 6 +++ 3 files changed, 32 insertions(+), 13 deletions(-) diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index efbd2c6e0..5118decf4 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -1128,7 +1128,7 @@ public class ImplicitFunctionBindingTests { @EnableAutoConfiguration public static class SupplierAndProcessorConfiguration { Many> sink = Sinks.many().unicast().onBackpressureBuffer(); - + @Bean public Supplier>> supplier() { return () -> sink.asFlux().doOnNext(v -> { diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index 8616fe3d0..3d1b28f4e 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -85,26 +85,39 @@ public class StreamBridgeTests { public static void before() { System.clearProperty("spring.cloud.function.definition"); } - + @Test // see SCF-985 void ensurePassThruFunctionIsNotPrePostProcessed() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( - TestChannelBinderConfiguration.getCompleteConfiguration(EmptyConfiguration.class)) - .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) { - + TestChannelBinderConfiguration.getCompleteConfiguration(EmptyConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) { + StreamBridge streamBridge = context.getBean(StreamBridge.class); OutputDestination outputDestination = context.getBean(OutputDestination.class); Message message = CloudEventMessageBuilder.withData("foo") // cloud event - .setType("CustomType") - .setSource("CustomSource") - .setSubject("CustomSubject") - .build(); + .setType("CustomType").setSource("CustomSource").setSubject("CustomSubject").build(); streamBridge.send("fooDestination", message); - Message messageReceived = outputDestination.receive(1000, "fooDestination"); - assertThat(CloudEventMessageUtils.getSource(messageReceived)).isEqualTo(URI.create("CustomSource")); - assertThat(CloudEventMessageUtils.getSubject(messageReceived)).isEqualTo("CustomSubject"); - assertThat(CloudEventMessageUtils.getType(messageReceived)).isEqualTo("CustomType"); + Message messageReceived = outputDestination.receive(1000, "fooDestination"); + assertThat(CloudEventMessageUtils.getSource(messageReceived)).isEqualTo(URI.create("CustomSource")); + assertThat(CloudEventMessageUtils.getSubject(messageReceived)).isEqualTo("CustomSubject"); + assertThat(CloudEventMessageUtils.getType(messageReceived)).isEqualTo("CustomType"); + } + } + + @Test // + void validateWARNifSendingToINputBinding() { // GH-2563 + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SimpleConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=echo", + "--spring.cloud.stream.function.bindings.echo-in-0=input")) { + + StreamBridge streamBridge = context.getBean(StreamBridge.class); + OutputDestination outputDestination = context.getBean(OutputDestination.class); + + streamBridge.send("input", "foo"); + // just validate that there is a WARN message about sending to the input binding } } diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index ee6363b42..e3dbe1c16 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -55,6 +55,7 @@ import org.springframework.messaging.support.GenericMessage; import org.springframework.util.Assert; import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; +import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -220,6 +221,11 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi if (messageChannel == null) { if (this.applicationContext.containsBean(destinationName)) { messageChannel = this.applicationContext.getBean(destinationName, MessageChannel.class); + String[] consumerBindingNames = this.bindingService.getConsumerBindingNames(); + if (ObjectUtils.containsElement(consumerBindingNames, destinationName)) { //GH-2563 + logger.warn("You seem to be sending data to the input binding. It is not " + + "recommended, since you are bypassing the binder and this the messaging system exposed by the binder."); + } } else { messageChannel = new DirectWithAttributesChannel();