diff --git a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/Application.java b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/Application.java index 19ff0dd462..b4562bfce3 100644 --- a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/Application.java +++ b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/Application.java @@ -18,11 +18,7 @@ package com.example.fraud; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.messaging.Source; -import org.springframework.context.annotation.Configuration; -@Configuration @SpringBootApplication public class Application { diff --git a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagePoller.java b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagePoller.java index 6a366b743e..676ea7d0a7 100644 --- a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagePoller.java +++ b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagePoller.java @@ -16,20 +16,27 @@ package com.example.fraud; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import reactor.core.publisher.EmitterProcessor; + import org.springframework.cloud.stream.messaging.Source; import org.springframework.integration.support.MessageBuilder; import org.springframework.stereotype.Component; @Component class MessagePoller { - private final Source source; - MessagePoller(Source source) { - this.source = source; + private static final Logger log = LoggerFactory.getLogger(MessagePoller.class); + + private final EmitterProcessor emitterProcessor; + + MessagePoller(EmitterProcessor emitterProcessor) { + this.emitterProcessor = emitterProcessor; } public void poll() { - this.source.output().send(MessageBuilder - .withPayload("{\"id\":\"99\",\"temperature\":\"123.45\"}").build()); + log.info("Emitting the message"); + this.emitterProcessor.onNext("{\"id\":\"99\",\"temperature\":\"123.45\"}"); } } diff --git a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessor.java b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagingConfiguration.java similarity index 58% rename from samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessor.java rename to samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagingConfiguration.java index 0a0b44a664..25b87065b8 100644 --- a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessor.java +++ b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MessagingConfiguration.java @@ -16,15 +16,23 @@ package com.example.fraud; -import org.springframework.cloud.stream.annotation.Output; -import org.springframework.cloud.stream.messaging.Sink; -import org.springframework.messaging.MessageChannel; +import java.util.function.Supplier; -interface MyProcessor extends Sink { +import reactor.core.publisher.EmitterProcessor; - String MY_OUTPUT = "my_output"; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; - @Output("my_output") - MessageChannel output(); +@Configuration +class MessagingConfiguration { + @Bean + EmitterProcessor sensorDataEmitter() { + return EmitterProcessor.create(); + } + + @Bean(name = "sensor-data") + Supplier sensorData() { + return () -> "{\"id\":\"99\",\"temperature\":\"123.45\"}"; + } } diff --git a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessorListener.java b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessorListener.java index 8c02597b04..f8417b76f9 100644 --- a/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessorListener.java +++ b/samples/standalone/dsl/http-server/src/main/java/com/example/fraud/MyProcessorListener.java @@ -22,6 +22,7 @@ import java.net.URISyntaxException; import java.net.URL; import java.nio.file.Files; import java.util.Arrays; +import java.util.function.Function; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -31,19 +32,16 @@ import org.springframework.cloud.stream.messaging.Sink; import org.springframework.integration.support.MessageBuilder; import org.springframework.stereotype.Component; -@Component -class MyProcessorListener { +@Component("my_output") +class MyProcessorListener implements Function { private static final Logger log = LoggerFactory.getLogger(MyProcessorListener.class); - private final MyProcessor processor; - private final byte[] expectedInput; private final byte[] expectedOutput; - MyProcessorListener(MyProcessor processor) { - this.processor = processor; + MyProcessorListener() { this.expectedInput = forFile("/contracts/messaging/input.pdf"); this.expectedOutput = forFile("/contracts/messaging/output.pdf"); } @@ -58,15 +56,13 @@ class MyProcessorListener { } } - @StreamListener(Sink.INPUT) - void listen(byte[] payload) { + @Override + public byte[] apply(byte[] payload) { log.info("Got the message!"); if (!Arrays.equals(payload, this.expectedInput)) { log.error("Input payload size is [" + payload.length + "] and the expected one is [" + this.expectedInput.length + "]"); throw new IllegalStateException("Wrong input"); } - this.processor.output() - .send(MessageBuilder.withPayload(this.expectedOutput).build()); + return this.expectedOutput; } - } diff --git a/samples/standalone/dsl/http-server/src/main/resources/application.properties b/samples/standalone/dsl/http-server/src/main/resources/application.properties index b00305698d..9f5532b835 100644 --- a/samples/standalone/dsl/http-server/src/main/resources/application.properties +++ b/samples/standalone/dsl/http-server/src/main/resources/application.properties @@ -1,6 +1,7 @@ -spring.cloud.stream.bindings.output.contentType=application/json -spring.cloud.stream.bindings.input.contentType=application/octet-stream -spring.cloud.stream.bindings.input.destination=bytes_input -spring.cloud.stream.bindings.my_output.contentType=application/octet-stream -spring.cloud.stream.bindings.my_output.destination=bytes_output +spring.cloud.function.definition=my_output;sensor-data +spring.cloud.stream.bindings.sensor-data-out-0.contentType=application/json +spring.cloud.stream.bindings.my_output-in-0.contentType=application/octet-stream +spring.cloud.stream.bindings.my_output-in-0.destination=bytes_input +spring.cloud.stream.bindings.my_output-out-0.contentType=application/octet-stream +spring.cloud.stream.bindings.my_output-out-0.destination=bytes_output server.port=0 \ No newline at end of file diff --git a/samples/standalone/dsl/http-server/src/test/java/com/example/fraud/MessagingBase.java b/samples/standalone/dsl/http-server/src/test/java/com/example/fraud/MessagingBase.java index c8c1845b8a..e701120aee 100644 --- a/samples/standalone/dsl/http-server/src/test/java/com/example/fraud/MessagingBase.java +++ b/samples/standalone/dsl/http-server/src/test/java/com/example/fraud/MessagingBase.java @@ -21,8 +21,11 @@ import org.junit.Before; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.ImportAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.contract.verifier.messaging.boot.AutoConfigureMessageVerifier; +import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.context.annotation.Configuration; import org.springframework.test.context.junit4.SpringRunner; import org.springframework.web.context.WebApplicationContext; @@ -34,7 +37,7 @@ import org.springframework.web.context.WebApplicationContext; * @author Marius Bogoevici */ @RunWith(SpringRunner.class) -@SpringBootTest(classes = Application.class, properties = "spring.cloud.stream.bindings.output.destination=sensor-data") +@SpringBootTest(classes = {MessagingBase.Config.class, Application.class}) @AutoConfigureMessageVerifier public abstract class MessagingBase { @@ -53,4 +56,10 @@ public abstract class MessagingBase { poller.poll(); } + @Configuration + @ImportAutoConfiguration(TestChannelBinderConfiguration.class) + static class Config { + + } + }