diff --git a/docs/src/main/asciidoc/spring-cloud-stream.adoc b/docs/src/main/asciidoc/spring-cloud-stream.adoc index fec7febdc..2089aa6c0 100644 --- a/docs/src/main/asciidoc/spring-cloud-stream.adoc +++ b/docs/src/main/asciidoc/spring-cloud-stream.adoc @@ -873,13 +873,13 @@ public void testMultipleFunctions() { OutputDestination outputDestination = context.getBean(OutputDestination.class); Message inputMessage = MessageBuilder.withPayload("Hello".getBytes()).build(); - inputDestination.send(inputMessage, 0); - inputDestination.send(inputMessage, 1); + inputDestination.send(inputMessage, "uppercase-in-0"); + inputDestination.send(inputMessage, "reverse-in-0"); - Message outputMessage = outputDestination.receive(0, 0); + Message outputMessage = outputDestination.receive(0, "uppercase-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("HELLO".getBytes()); - outputMessage = outputDestination.receive(0, 1); + outputMessage = outputDestination.receive(0, "uppercase-out-1"); assertThat(outputMessage.getPayload()).isEqualTo("olleH".getBytes()); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java index 86e5b787a..12040958a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java @@ -44,7 +44,7 @@ abstract class AbstractDestination { } SubscribableChannel getChannelByName(String name) { - //name = name + ".destination"; + name = name.endsWith(".destination") ? name : name + ".destination"; for (AbstractSubscribableChannel subscribableChannel : channels) { if (subscribableChannel.getBeanName().equals(name)) { return subscribableChannel; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/InputDestination.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/InputDestination.java index 1be33bd2d..8351eb57d 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/InputDestination.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/InputDestination.java @@ -29,18 +29,30 @@ import org.springframework.messaging.Message; public class InputDestination extends AbstractDestination { /** - * Allows the {@link Message} to be sent to a Binder to be delegated to binder's input - * destination (e.g., Processor.INPUT). + * Allows the {@link Message} to be sent to a Binder to be delegated to a default binding + * destination (e.g., "function-in-0" for cases where you only have a single function with the name 'function'). * @param message message to send */ public void send(Message message) { this.getChannel(0).send(message); } + /** + * @param message message to send + * @param inputIndex input index + * @deprecated since 3.0.2 in favor of {@link #receive(long, String)} where you should use the actual binding name (e.g., "foo-in-0") + */ + @Deprecated public void send(Message message, int inputIndex) { this.getChannel(inputIndex).send(message); } + /** + * Allows the {@link Message} to be sent to a Binder to be delegated to a named binding + * destination (e.g., "function-in-0" for cases where you want to send input to function with the name 'function'). + * @param message message to send + * @param bindingName binding name + */ public void send(Message message, String bindingName) { this.getChannelByName(bindingName).send(message); } 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 efccab0a3..99da07381 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 @@ -18,7 +18,6 @@ package org.springframework.cloud.stream.binder.test; import java.util.ArrayList; import java.util.LinkedHashMap; -import java.util.List; import java.util.Map; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedTransferQueue; @@ -40,6 +39,7 @@ public class OutputDestination extends AbstractDestination { public Message receive(long timeout, String bindingName) { try { + bindingName = bindingName.endsWith(".destination") ? bindingName : bindingName + ".destination"; return this.messageQueues.get(bindingName).poll(timeout, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { @@ -51,7 +51,9 @@ public class OutputDestination extends AbstractDestination { * Allows to access {@link Message}s received by this {@link OutputDestination}. * @param timeout how long to wait before giving up * @return received message + * @deprecated since 3.0.2 in favor of {@link #receive(long, String)} where you should use the actual binding name (e.g., "foo-in-0") */ + @Deprecated public Message receive(long timeout, int bindingIndex) { try { BlockingQueue> destinationQueue = (new ArrayList<>(this.messageQueues.values())).get(bindingIndex); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java index 5b1e782aa..7cf50e4af 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java @@ -141,12 +141,12 @@ public class MultipleInputOutputFunctionTests { Message stringInputMessage = MessageBuilder.withPayload("one".getBytes()).build(); Message integerInputMessage = MessageBuilder.withPayload("1".getBytes()).build(); - inputDestination.send(stringInputMessage, 0); - inputDestination.send(integerInputMessage, 1); + inputDestination.send(stringInputMessage, "multiInputSingleOutput-in-0"); + inputDestination.send(integerInputMessage, "multiInputSingleOutput-in-1"); Message outputMessage = outputDestination.receive(); assertThat(outputMessage.getPayload()).isEqualTo("one".getBytes()); - outputMessage = outputDestination.receive(); + outputMessage = outputDestination.receive(0, "multiInputSingleOutput-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("1".getBytes()); } } @@ -165,14 +165,14 @@ public class MultipleInputOutputFunctionTests { OutputDestination outputDestination = context.getBean(OutputDestination.class); for (int i = 0; i < 10; i++) { - inputDestination.send(MessageBuilder.withPayload(String.valueOf(i).getBytes()).build()); + inputDestination.send(MessageBuilder.withPayload(String.valueOf(i).getBytes()).build(), "singleInputMultipleOutputs-in-0"); } int counter = 0; for (int i = 0; i < 5; i++) { - Message even = outputDestination.receive(0, 0); + Message even = outputDestination.receive(0, "singleInputMultipleOutputs-out-0"); assertThat(even.getPayload()).isEqualTo(("EVEN: " + String.valueOf(counter++)).getBytes()); - Message odd = outputDestination.receive(0, 1); + Message odd = outputDestination.receive(0, "singleInputMultipleOutputs-out-1"); assertThat(odd.getPayload()).isEqualTo(("ODD: " + String.valueOf(counter++)).getBytes()); } } @@ -192,13 +192,13 @@ public class MultipleInputOutputFunctionTests { OutputDestination outputDestination = context.getBean(OutputDestination.class); Message inputMessage = MessageBuilder.withPayload("Hello".getBytes()).build(); - inputDestination.send(inputMessage, 0); - inputDestination.send(inputMessage, 1); + inputDestination.send(inputMessage, "uppercase-in-0"); + inputDestination.send(inputMessage, "reverse-in-0"); - Message outputMessage = outputDestination.receive(0, 0); + Message outputMessage = outputDestination.receive(0, "uppercase-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("HELLO".getBytes()); - outputMessage = outputDestination.receive(0, 1); + outputMessage = outputDestination.receive(0, "reverse-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("olleH".getBytes()); } } @@ -217,13 +217,13 @@ public class MultipleInputOutputFunctionTests { OutputDestination outputDestination = context.getBean(OutputDestination.class); Message inputMessage = MessageBuilder.withPayload("Hello".getBytes()).build(); - inputDestination.send(inputMessage, 0); - inputDestination.send(inputMessage, 1); + inputDestination.send(inputMessage, "uppercasereverse-in-0"); + inputDestination.send(inputMessage, "reverseuppercase-in-0"); - Message outputMessage = outputDestination.receive(0, 0); + Message outputMessage = outputDestination.receive(0, "uppercasereverse-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("OLLEH".getBytes()); - outputMessage = outputDestination.receive(0, 1); + outputMessage = outputDestination.receive(0, "reverseuppercase-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("OLLEH".getBytes()); } } @@ -245,12 +245,12 @@ public class MultipleInputOutputFunctionTests { Message stringInputMessage = MessageBuilder.withPayload("ricky".getBytes()).build(); Message integerInputMessage = MessageBuilder.withPayload("bobby".getBytes()).build(); - inputDestination.send(stringInputMessage, 0); - inputDestination.send(integerInputMessage, 1); + inputDestination.send(stringInputMessage, "multiInputSingleOutput-in-0"); + inputDestination.send(integerInputMessage, "multiInputSingleOutput-in-1"); - Message outputMessage = outputDestination.receive(); + Message outputMessage = outputDestination.receive(1000, "multiInputSingleOutput-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("RICKY".getBytes()); - outputMessage = outputDestination.receive(); + outputMessage = outputDestination.receive(1000, "multiInputSingleOutput-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("BOBBY".getBytes()); } } @@ -273,13 +273,11 @@ public class MultipleInputOutputFunctionTests { Message stringInputMessage = MessageBuilder.withPayload("ricky".getBytes()).build(); Message integerInputMessage = MessageBuilder.withPayload("bobby".getBytes()).build(); - inputDestination.send(stringInputMessage, 0); - inputDestination.send(integerInputMessage, 1); + inputDestination.send(stringInputMessage, "multiInputSingleOutput2-in-0"); + inputDestination.send(integerInputMessage, "multiInputSingleOutput2-in-1"); - Message outputMessage = outputDestination.receive(2500); + Message outputMessage = outputDestination.receive(1000, "multiInputSingleOutput2-out-0"); assertThat(outputMessage.getPayload()).isEqualTo("rickybobby".getBytes()); -// outputMessage = outputDestination.receive(); -// assertThat(outputMessage.getPayload()).isEqualTo("BOBBY".getBytes()); } } @@ -384,6 +382,7 @@ public class MultipleInputOutputFunctionTests { return Person.class.isAssignableFrom(clazz); } + @Override @Nullable protected Object convertFromInternal( Message message, Class targetClass, @Nullable Object conversionHint) { @@ -392,6 +391,7 @@ public class MultipleInputOutputFunctionTests { return person; } + @Override @Nullable protected Object convertToInternal( Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) { @@ -410,6 +410,7 @@ public class MultipleInputOutputFunctionTests { return Employee.class.isAssignableFrom(clazz); } + @Override @Nullable protected Object convertFromInternal( Message message, Class targetClass, @Nullable Object conversionHint) { @@ -418,6 +419,7 @@ public class MultipleInputOutputFunctionTests { return person; } + @Override @Nullable protected Object convertToInternal( Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) {