From f2b9d3ecdc2f5690ec4cdff5896863eaf1b4cfdf Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 5 May 2021 20:33:35 +0200 Subject: [PATCH] GH-2152 Add support for honoring timeout for cases where destination may not yet exist Resolves #2152 polishing --- .../stream/binder/test/OutputDestination.java | 29 ++++++++++----- .../ImplicitFunctionBindingTests.java | 14 +++---- .../SourceToFunctionsSupportTests.java | 37 +++---------------- .../stream/function/StreamBridgeTests.java | 24 ++++++++++++ 4 files changed, 56 insertions(+), 48 deletions(-) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java index 5f73b083f..a86a34137 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java @@ -1,5 +1,5 @@ /* - * Copyright 2017-2019 the original author or authors. + * Copyright 2017-2021 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. @@ -17,12 +17,15 @@ package org.springframework.cloud.stream.binder.test; import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.Map; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedTransferQueue; import java.util.concurrent.TimeUnit; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +import org.springframework.integration.channel.AbstractSubscribableChannel; import org.springframework.messaging.Message; import org.springframework.util.StringUtils; @@ -36,12 +39,14 @@ import org.springframework.util.StringUtils; */ public class OutputDestination extends AbstractDestination { - private final Map>> messageQueues = new LinkedHashMap<>(); + private final Log log = LogFactory.getLog(OutputDestination.class); + + private final ConcurrentHashMap>> messageQueues = new ConcurrentHashMap<>(); public Message receive(long timeout, String bindingName) { try { bindingName = bindingName.endsWith(".destination") ? bindingName : bindingName + ".destination"; - return this.messageQueues.get(bindingName).poll(timeout, TimeUnit.MILLISECONDS); + return this.outputQueue(bindingName).poll(timeout, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -81,11 +86,13 @@ public class OutputDestination extends AbstractDestination { */ @Deprecated public Message receive(long timeout, int bindingIndex) { + log.warn("!!!While 'receive(long timeout, int bindingIndex)' method may still work it is deprecated no longer supported. " + + "It will be removed after 3.1.3 release. Please use 'receive(long timeout, String bindingName)'"); try { BlockingQueue> destinationQueue = (new ArrayList<>(this.messageQueues.values())).get(bindingIndex); return destinationQueue.poll(timeout, TimeUnit.MILLISECONDS); } - catch (InterruptedException e) { + catch (Exception e) { Thread.currentThread().interrupt(); } return null; @@ -106,11 +113,13 @@ public class OutputDestination extends AbstractDestination { @SuppressWarnings("unchecked") @Override void afterChannelIsSet(int channelIndex, String bindingName) { - if (!this.messageQueues.containsKey(bindingName)) { - BlockingQueue> messageQueue = new LinkedTransferQueue<>(); - this.messageQueues.put(bindingName, messageQueue); - this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message) message)); + if (((AbstractSubscribableChannel) this.getChannelByName(bindingName)).getSubscriberCount() < 1) { + this.getChannelByName(bindingName).subscribe(message -> this.outputQueue(bindingName).offer((Message) message)); } } + private BlockingQueue> outputQueue(String bindingName) { + this.messageQueues.putIfAbsent(bindingName, new LinkedTransferQueue<>()); + return this.messageQueues.get(bindingName); + } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 4a22519a0..f6cb11ac9 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -286,14 +286,14 @@ public class ImplicitFunctionBindingTests { inputDestination.send(inputMessage); inputDestination.send(inputMessage); inputDestination.send(inputMessage); - assertThat(new String(outputDestination.receive(2000).getPayload())).isEqualTo("HelloHelloHello"); - assertThat(new String(outputDestination.receive(2000).getPayload())).isEqualTo(""); + assertThat(new String(outputDestination.receive(2000, "aggregate-out-0").getPayload())).isEqualTo("HelloHelloHello"); + assertThat(new String(outputDestination.receive(2000, "aggregate-out-0").getPayload())).isEqualTo(""); inputDestination.send(inputMessage); inputDestination.send(inputMessage); inputDestination.send(inputMessage); inputDestination.send(inputMessage); - assertThat(new String(outputDestination.receive(2000).getPayload())).isEqualTo("HelloHelloHelloHello"); + assertThat(new String(outputDestination.receive(2000, "aggregate-out-0").getPayload())).isEqualTo("HelloHelloHelloHello"); } } @@ -565,16 +565,16 @@ public class ImplicitFunctionBindingTests { OutputDestination outputDestination = context.getBean(OutputDestination.class); - Message outputMessage = outputDestination.receive(2000); + Message outputMessage = outputDestination.receive(2000, "supplierfunctionA-out-0"); Long value = Long.parseLong(new String(outputMessage.getPayload())); - outputMessage = outputDestination.receive(5000); + outputMessage = outputDestination.receive(5000, "supplierfunctionA-out-0"); assertThat(Long.parseLong(new String(outputMessage.getPayload())) - value).isGreaterThanOrEqualTo(1000); - outputMessage = outputDestination.receive(5000); + outputMessage = outputDestination.receive(5000, "supplierfunctionA-out-0"); assertThat(Long.parseLong(new String(outputMessage.getPayload())) - value).isGreaterThanOrEqualTo(1000); - outputMessage = outputDestination.receive(5000); + outputMessage = outputDestination.receive(5000, "supplierfunctionA-out-0"); assertThat(Long.parseLong(new String(outputMessage.getPayload())) - value).isGreaterThanOrEqualTo(1000); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java index cf9fe2273..fdd8729d2 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java @@ -240,40 +240,15 @@ public class SourceToFunctionsSupportTests { OutputDestination target = context.getBean(OutputDestination.class); - assertThat(new String(target.receive(2000).getPayload())).isEqualTo("1"); - assertThat(new String(target.receive(2000).getPayload())).isEqualTo("2"); - assertThat(new String(target.receive(2000).getPayload())).isEqualTo("3"); - assertThat(new String(target.receive(2000).getPayload())).isEqualTo("4"); - assertThat(new String(target.receive(2000).getPayload())).isEqualTo("5"); - assertThat(new String(target.receive(2000).getPayload())).isEqualTo("6"); - - //assertThat(context.getBean("supplierInitializer")).isNotEqualTo(null); + assertThat(new String(target.receive(2000, "simpleStreamSupplier-out-0").getPayload())).isEqualTo("1"); + assertThat(new String(target.receive(2000, "simpleStreamSupplier-out-0").getPayload())).isEqualTo("2"); + assertThat(new String(target.receive(2000, "simpleStreamSupplier-out-0").getPayload())).isEqualTo("3"); + assertThat(new String(target.receive(2000, "simpleStreamSupplier-out-0").getPayload())).isEqualTo("4"); + assertThat(new String(target.receive(2000, "simpleStreamSupplier-out-0").getPayload())).isEqualTo("5"); + assertThat(new String(target.receive(2000, "simpleStreamSupplier-out-0").getPayload())).isEqualTo("6"); } } - @Test - public void testMultipleSuppliers() { - try (ConfigurableApplicationContext context = new SpringApplicationBuilder( - TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class, - MultipleSupplierConfiguration.class)).web(WebApplicationType.NONE).run( - "--spring.cloud.function.definition=supplier1;supplier2", - "--spring.jmx.enabled=false" -// "--spring.cloud.stream.function.bindings.supplier1-out-0=output1", -// "--spring.cloud.stream.function.bindings.supplier2-out-0=output2" - )) { - - OutputDestination target = context.getBean(OutputDestination.class); - -// assertThat(new String(target.receive(2000).getPayload())).isEqualTo("1"); -// assertThat(new String(target.receive(2000).getPayload())).isEqualTo("2"); -// assertThat(new String(target.receive(2000).getPayload())).isEqualTo("3"); -// assertThat(new String(target.receive(2000).getPayload())).isEqualTo("4"); -// assertThat(new String(target.receive(2000).getPayload())).isEqualTo("5"); -// assertThat(new String(target.receive(2000).getPayload())).isEqualTo("6"); -// -// assertThat(context.getBean("supplierInitializer")).isNotEqualTo(null); - } - } @EnableAutoConfiguration public static class MessageFluxSupplierConfiguration { diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index 1ecdefe18..b04215218 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.function; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; import java.util.function.Function; @@ -57,6 +60,27 @@ public class StreamBridgeTests { System.clearProperty("spring.cloud.function.definition"); } + @Test + public void testDelayedSend() { + ScheduledExecutorService executor = Executors.newScheduledThreadPool(1); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(ConsumerConfiguration.class, EmptyConfiguration.class)) + .web(WebApplicationType.NONE).run( + "--spring.jmx.enabled=false")) { + + StreamBridge bridge = context.getBean(StreamBridge.class); + executor.schedule(() -> bridge.send("blah", "hello foo"), 5000, TimeUnit.MILLISECONDS); + + OutputDestination outputDestination = context.getBean(OutputDestination.class); + Message message = outputDestination.receive(10000, "blah"); + assertThat(message).isNotNull(); + assertThat(new String(message.getPayload())).isEqualTo("hello foo"); + } + finally { + executor.shutdownNow(); + } + } + @Test public void testWithInterceptor() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration