diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index b06d65cb2..bcd106ba1 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -414,8 +414,8 @@ Not supported when the `containerType` is `direct`. + Default: `1`. -[[rabbitmq-stream-consumer]] -=== Initial Consumer Support for the RabbitMQ Stream Plugin +[[rabbitmq-stream]] +=== Initial Support for the RabbitMQ Stream Plugin Basic support for the https://rabbitmq.com/stream.html[RabbitMQ Stream Plugin] is now provided. To enable this feature, you must add the `spring-rabbit-stream` jar to the class path - it must be the same version as `spring-amqp` and `spring-rabbit`. @@ -1027,7 +1027,6 @@ public class Application { catch (ExecutionException | TimeoutException e) { throw new IllegalStateException(e); } - }); }; } @@ -1210,38 +1209,3 @@ For negatively acknowledged confirmations, the payload is a `NackedAmqpMessageEx There is no automatic handling of these exceptions (such as sending to a <>). You can consume these exceptions with your own Spring Integration flow. - -[[rabbitmq-stream-producer]] -=== Initial Producer Support for the RabbitMQ Stream Plugin - -Basic support for the https://rabbitmq.com/stream.html[RabbitMQ Stream Plugin] is now provided. -To enable this feature, you must add the `spring-rabbit-stream` jar to the class path - it must be the same version as `spring-amqp` and `spring-rabbit`. - -IMPORTANT: The producer properties described above are not supported when you set the `producerType` property to `STREAM_SYNC` or `STREAM_ASYNC`. - -To configure the binder to use a stream `ProducerType`, you must add an `Environment` `@Bean` and, optionally, a customizer to customize the message handler. - -==== -[source, java] ----- -@Bean -Environment streamEnv() { - return Environment.builder() - .build(); -} - -@Bean -ProducerMessageHandlerCustomizer handlerCustomizer() { - return (hand, dest) -> { - RabbitStreamMessageHandler handler = (RabbitStreamMessageHandler) hand; - handler.setConfirmTimeout(5000); - ((RabbitStreamTemplate) handler.getStreamOperations()).setProducerCustomizer( - (name, builder) -> { - ... - }); - }; -} ----- -==== - -Refer to the https://rabbitmq.github.io/rabbitmq-stream-java-client/stable/htmlsingle/[RabbitMQ Stream Java Client documentation] for information about configuring the environment and producer builder. diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java index b72dcd91d..0dd3acd30 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java @@ -231,13 +231,6 @@ public abstract class RabbitCommonProperties { */ private boolean dlqSingleActiveConsumer; - /** - * The bean name of a stream message converter to convert from a Spring AMQP Message - * to a Stream Message. - * @since 3.2 - */ - private String streamStreamMessageConverterBeanName; - public String getExchangeType() { return this.exchangeType; } @@ -543,14 +536,6 @@ public abstract class RabbitCommonProperties { this.dlqSingleActiveConsumer = dlqSingleActiveConsumer; } - public String getStreamStreamMessageConverterBeanName() { - return this.streamStreamMessageConverterBeanName; - } - - public void setStreamStreamMessageConverterBeanName(String streamStreamMessageConverterBeanName) { - this.streamStreamMessageConverterBeanName = streamStreamMessageConverterBeanName; - } - public static class QuorumConfig { private boolean enabled; diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java index bbfd7f64b..629f8e2e3 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java @@ -20,7 +20,6 @@ import javax.validation.constraints.Min; import org.springframework.amqp.core.MessageDeliveryMode; import org.springframework.expression.Expression; -import org.springframework.util.Assert; /** * @author Marius Bogoevici @@ -28,28 +27,6 @@ import org.springframework.util.Assert; */ public class RabbitProducerProperties extends RabbitCommonProperties { - /** - * Determines the producer type. - * @since 3.2 - */ - public enum ProducerType { - - /** - * RabbitMQ Stream producer - blocks until confirm received. - */ - STREAM_SYNC, - - /** - * RabbitMQ Stream producer - does not block. - */ - STREAM_ASYNC, - - /** - * Classic AMQP producer. - */ - AMQP - } - /** * true to compress messages. */ @@ -124,20 +101,6 @@ public class RabbitProducerProperties extends RabbitCommonProperties { */ private boolean useConfirmHeader; - /** - * When STREAM_SYNC or STREAM_ASYNC, create a RabbitMQ Stream producer instead of an - * AMQP producer. - * @since 3.2 - */ - private ProducerType producerType = ProducerType.AMQP; - - /** - * The bean name of a message converter to convert from spring-messaging Message to - * a Spring AMQP Message. - * @since 3.2 - */ - private String streamMessageConverterBeanName; - /** * @deprecated - use {@link #setHeaderPatterns(String[])}. * @param requestHeaderPatterns the patterns. @@ -263,21 +226,4 @@ public class RabbitProducerProperties extends RabbitCommonProperties { this.useConfirmHeader = useConfirmHeader; } - public ProducerType getProducerType() { - return this.producerType; - } - - public void setProducerType(ProducerType producerType) { - Assert.notNull(producerType, "'producerType' cannot be null"); - this.producerType = producerType; - } - - public String getStreamMessageConverterBeanName() { - return this.streamMessageConverterBeanName; - } - - public void setStreamMessageConverterBeanName(String streamMessageConverterBeanName) { - this.streamMessageConverterBeanName = streamMessageConverterBeanName; - } - } diff --git a/spring-cloud-stream-binder-rabbit/pom.xml b/spring-cloud-stream-binder-rabbit/pom.xml index 235c3f3cb..3e4de3806 100644 --- a/spring-cloud-stream-binder-rabbit/pom.xml +++ b/spring-cloud-stream-binder-rabbit/pom.xml @@ -90,12 +90,6 @@ spring-cloud-stream-binder-rabbit-test-support test - - org.testcontainers - rabbitmq - 1.15.3 - test - org.apache.httpcomponents diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index 98197b58c..344334763 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -81,7 +81,6 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerP import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties.ContainerType; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedBindingProperties; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties.ProducerType; import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.config.MessageSourceCustomizer; @@ -99,7 +98,6 @@ import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter.BatchMode; import org.springframework.integration.amqp.inbound.AmqpMessageSource; import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; -import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.amqp.support.AmqpMessageHeaderErrorMessageStrategy; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; import org.springframework.integration.channel.AbstractMessageChannel; @@ -310,21 +308,6 @@ public class RabbitMessageChannelBinder extends String destination = StringUtils.isEmpty(prefix) ? exchangeName : exchangeName.substring(prefix.length()); RabbitProducerProperties extendedProperties = producerProperties.getExtension(); - final MessageHandler endpoint; - if (!ProducerType.AMQP.equals(producerProperties.getExtension().getProducerType())) { - endpoint = StreamUtils.createStreamMessageHandler(producerDestination, producerProperties, errorChannel, - destination, extendedProperties, getApplicationContext(), this::configureHeaderMapper); - } - else { - endpoint = amqpHandler(producerDestination, producerProperties, errorChannel, - destination, extendedProperties); - } - return endpoint; - } - - private AmqpOutboundEndpoint amqpHandler(final ProducerDestination producerDestination, - ExtendedProducerProperties producerProperties, MessageChannel errorChannel, - String destination, RabbitProducerProperties extendedProperties) { final AmqpOutboundEndpoint endpoint = new AmqpOutboundEndpoint( buildRabbitTemplate(extendedProperties, errorChannel != null || extendedProperties.isUseConfirmHeader())); @@ -374,7 +357,16 @@ public class RabbitMessageChannelBinder extends endpoint.setDelayExpression(extendedProperties.getDelayExpression()); } } - endpoint.setHeaderMapper(configureHeaderMapper(extendedProperties)); + DefaultAmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.outboundMapper(); + List headerPatterns = new ArrayList<>(extendedProperties.getHeaderPatterns().length + 3); + headerPatterns.add("!" + BinderHeaders.PARTITION_HEADER); + headerPatterns.add("!" + IntegrationMessageHeaderAccessor.SOURCE_DATA); + headerPatterns.add("!" + IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT); + headerPatterns.add("!rabbitmq_streamContext"); + headerPatterns.addAll(Arrays.asList(extendedProperties.getHeaderPatterns())); + mapper.setRequestHeaderNames( + headerPatterns.toArray(new String[headerPatterns.size()])); + endpoint.setHeaderMapper(mapper); endpoint.setDefaultDeliveryMode(extendedProperties.getDeliveryMode()); endpoint.setBeanFactory(this.getBeanFactory()); if (errorChannel != null) { @@ -405,19 +397,6 @@ public class RabbitMessageChannelBinder extends return endpoint; } - private AmqpHeaderMapper configureHeaderMapper(RabbitProducerProperties extendedProperties) { - DefaultAmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.outboundMapper(); - List headerPatterns = new ArrayList<>(extendedProperties.getHeaderPatterns().length + 3); - headerPatterns.add("!" + BinderHeaders.PARTITION_HEADER); - headerPatterns.add("!" + IntegrationMessageHeaderAccessor.SOURCE_DATA); - headerPatterns.add("!" + IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT); - headerPatterns.add("!rabbitmq_streamContext"); - headerPatterns.addAll(Arrays.asList(extendedProperties.getHeaderPatterns())); - mapper.setRequestHeaderNames( - headerPatterns.toArray(new String[headerPatterns.size()])); - return mapper; - } - @Override protected void postProcessOutputChannel(MessageChannel outputChannel, ExtendedProducerProperties producerProperties) { @@ -514,7 +493,7 @@ public class RabbitMessageChannelBinder extends adapter.setBatchMode(BatchMode.EXTRACT_PAYLOADS_WITH_HEADERS); } if (extension.getContainerType().equals(ContainerType.STREAM)) { - StreamUtils.configureAdapter(adapter); + StreamContainerUtils.configureAdapter(adapter); } return adapter; } @@ -524,7 +503,7 @@ public class RabbitMessageChannelBinder extends RabbitConsumerProperties extension) { if (extension.getContainerType().equals(ContainerType.STREAM)) { - return StreamUtils.createContainer(consumerDestination, group, properties, destination, extension, + return StreamContainerUtils.createContainer(consumerDestination, group, properties, destination, extension, getApplicationContext()); } boolean directContainer = extension.getContainerType() diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitStreamMessageHandler.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitStreamMessageHandler.java deleted file mode 100644 index 6930d42be..000000000 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitStreamMessageHandler.java +++ /dev/null @@ -1,259 +0,0 @@ -/* - * Copyright 2021-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. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.binder.rabbit; - -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; - -import org.springframework.amqp.core.MessageProperties; -import org.springframework.amqp.support.AmqpHeaders; -import org.springframework.amqp.support.converter.ContentTypeDelegatingMessageConverter; -import org.springframework.amqp.support.converter.MessageConverter; -import org.springframework.context.Lifecycle; -import org.springframework.integration.amqp.support.AmqpHeaderMapper; -import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; -import org.springframework.integration.handler.AbstractMessageHandler; -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandler; -import org.springframework.messaging.MessageHandlingException; -import org.springframework.messaging.MessageHeaders; -import org.springframework.rabbit.stream.producer.RabbitStreamOperations; -import org.springframework.rabbit.stream.support.StreamMessageProperties; -import org.springframework.util.Assert; -import org.springframework.util.MimeType; -import org.springframework.util.concurrent.ListenableFuture; -import org.springframework.util.concurrent.SuccessCallback; - -/** - * {@link MessageHandler} based on {@link RabbitStreamOperations}. - * - * TODO: This class will move to Spring Integration in 6.0. - * - * @author Gary Russell - * @since 3.2 - * - */ -public class RabbitStreamMessageHandler extends AbstractMessageHandler implements Lifecycle { - - private static final int DEFAULT_CONFIRM_TIMEOUT = 10_000; - - private final RabbitStreamOperations streamOperations; - - private boolean sync; - - private long confirmTimeout = DEFAULT_CONFIRM_TIMEOUT; - - private SuccessCallback> successCallback = msg -> { }; - - private FailureCallback failureCallback = (msg, ex) -> { }; - - private AmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.outboundMapper(); - - private boolean headersMappedLast; - - /** - * Create an instance with the provided {@link RabbitStreamOperations}. - * @param streamOperations the operations. - */ - public RabbitStreamMessageHandler(RabbitStreamOperations streamOperations) { - Assert.notNull(streamOperations, "'streamOperations' cannot be null"); - this.streamOperations = streamOperations; - } - - /** - * Set a callback to be invoked when a send is successful. - * @param successCallback the callback. - */ - public void setSuccessCallback(SuccessCallback> successCallback) { - Assert.notNull(successCallback, "'successCallback' cannot be null"); - this.successCallback = successCallback; - } - - /** - * Set a callback to be invoked when a send fails. - * @param failureCallback the callback. - */ - public void setFailureCallback(FailureCallback failureCallback) { - Assert.notNull(failureCallback, "'failureCallback' cannot be null"); - this.failureCallback = failureCallback; - } - - /** - * Set to true to wait for a confirmation. - * @param sync true to wait. - * @see #setConfirmTimeout(long) - */ - public void setSync(boolean sync) { - this.sync = sync; - } - - /** - * Set the confirm timeout. - * @param confirmTimeout the timeout. - * @see #setSync(boolean) - */ - public void setConfirmTimeout(long confirmTimeout) { - this.confirmTimeout = confirmTimeout; - } - - /** - * Set a custom {@link AmqpHeaderMapper} for mapping request and reply headers. - * Defaults to {@link DefaultAmqpHeaderMapper#outboundMapper()}. - * @param headerMapper the {@link AmqpHeaderMapper} to use. - */ - public void setHeaderMapper(AmqpHeaderMapper headerMapper) { - Assert.notNull(headerMapper, "headerMapper must not be null"); - this.headerMapper = headerMapper; - } - - /** - * When mapping headers for the outbound message, determine whether the headers are - * mapped before the message is converted, or afterwards. This only affects headers - * that might be added by the message converter. When false, the converter's headers - * win; when true, any headers added by the converter will be overridden (if the - * source message has a header that maps to those headers). You might wish to set this - * to true, for example, when using a - * {@link org.springframework.amqp.support.converter.SimpleMessageConverter} with a - * String payload that contains json; the converter will set the content type to - * {@code text/plain} which can be overridden to {@code application/json} by setting - * the {@link AmqpHeaders#CONTENT_TYPE} message header. Default: false. - * @param headersMappedLast true if headers are mapped after conversion. - */ - public void setHeadersMappedLast(boolean headersMappedLast) { - this.headersMappedLast = headersMappedLast; - } - - /** - * Return the {@link RabbitStreamOperations}. - * @return the operations. - */ - public RabbitStreamOperations getStreamOperations() { - return this.streamOperations; - } - - @Override - protected void handleMessageInternal(Message requestMessage) { - ListenableFuture future; - com.rabbitmq.stream.Message streamMessage; - if (requestMessage.getPayload() instanceof com.rabbitmq.stream.Message) { - streamMessage = (com.rabbitmq.stream.Message) requestMessage.getPayload(); - } - else { - MessageConverter converter = streamOperations.messageConverter(); - org.springframework.amqp.core.Message amqpMessage = mapMessage(requestMessage, converter, - this.headerMapper, this.headersMappedLast); - streamMessage = this.streamOperations.streamMessageConverter().fromMessage(amqpMessage); - } - future = this.streamOperations.send(streamMessage); - handleConfirms(requestMessage, future); - } - - private void handleConfirms(Message message, ListenableFuture future) { - future.addCallback(bool -> this.successCallback.onSuccess(message), - ex -> this.failureCallback.failure(message, ex)); - if (this.sync) { - try { - future.get(this.confirmTimeout, TimeUnit.MILLISECONDS); - } - catch (InterruptedException ex) { - Thread.currentThread().interrupt(); - throw new MessageHandlingException(message, ex); - } - catch (ExecutionException | TimeoutException ex) { - throw new MessageHandlingException(message, ex); - } - } - } - - /* - * TODO Copied/modified from MapppingUtils until SI 6.0 - */ - private static org.springframework.amqp.core.Message mapMessage(Message message, - MessageConverter converter, AmqpHeaderMapper headerMapper, boolean headersMappedLast) { - - MessageProperties amqpMessageProperties = new StreamMessageProperties(); - org.springframework.amqp.core.Message amqpMessage; - if (!headersMappedLast) { - mapHeaders(message.getHeaders(), amqpMessageProperties, headerMapper); - } - if (converter instanceof ContentTypeDelegatingMessageConverter && headersMappedLast) { - String contentType = contentTypeAsString(message.getHeaders()); - if (contentType != null) { - amqpMessageProperties.setContentType(contentType); - } - } - amqpMessage = converter.toMessage(message.getPayload(), amqpMessageProperties); - if (headersMappedLast) { - mapHeaders(message.getHeaders(), amqpMessageProperties, headerMapper); - } - return amqpMessage; - } - - private static void mapHeaders(MessageHeaders messageHeaders, MessageProperties amqpMessageProperties, - AmqpHeaderMapper headerMapper) { - - headerMapper.fromHeadersToRequest(messageHeaders, amqpMessageProperties); - } - - private static String contentTypeAsString(MessageHeaders headers) { - Object contentType = headers.get(AmqpHeaders.CONTENT_TYPE); - if (contentType instanceof MimeType) { - contentType = contentType.toString(); - } - if (contentType instanceof String) { - return (String) contentType; - } - else if (contentType != null) { - throw new IllegalArgumentException(AmqpHeaders.CONTENT_TYPE - + " header must be a MimeType or String, found: " + contentType.getClass().getName()); - } - return null; - } - /* - * End copied/modified from MappingUtils - */ - - @Override - public void start() { - } - - @Override - public void stop() { - this.streamOperations.close(); - } - - @Override - public boolean isRunning() { - return true; - } - - /** - * Callback for when publishing fails. - */ - public interface FailureCallback { - - /** - * Message publish failure. - * @param message the message. - * @param throwable the throwable. - */ - void failure(Message message, Throwable throwable); - - } - -} diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java new file mode 100644 index 000000000..1163a07e2 --- /dev/null +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java @@ -0,0 +1,257 @@ +/* + * Copyright 2021-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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.rabbit; + +import java.nio.charset.Charset; +import java.nio.charset.StandardCharsets; +import java.util.Map; +import java.util.UUID; +import java.util.function.Supplier; + +import com.rabbitmq.stream.Codec; +import com.rabbitmq.stream.Environment; +import com.rabbitmq.stream.MessageBuilder; +import com.rabbitmq.stream.MessageBuilder.ApplicationPropertiesBuilder; +import com.rabbitmq.stream.MessageBuilder.PropertiesBuilder; +import com.rabbitmq.stream.Properties; +import com.rabbitmq.stream.codec.WrapperMessageBuilder; + +import org.springframework.amqp.core.Message; +import org.springframework.amqp.core.MessageProperties; +import org.springframework.amqp.rabbit.listener.MessageListenerContainer; +import org.springframework.amqp.support.converter.MessageConversionException; +import org.springframework.amqp.utils.JavaUtils; +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.context.support.AbstractApplicationContext; +import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; +import org.springframework.integration.amqp.support.AmqpHeaderMapper; +import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; +import org.springframework.lang.Nullable; +import org.springframework.messaging.MessageHeaders; +import org.springframework.rabbit.stream.listener.ConsumerCustomizer; +import org.springframework.rabbit.stream.listener.StreamListenerContainer; +import org.springframework.rabbit.stream.support.StreamMessageProperties; +import org.springframework.rabbit.stream.support.converter.StreamMessageConverter; +import org.springframework.util.Assert; + +/** + * Utilities for stream containers. Used to prevent a hard runtime dependency on + * spring-rabbit-stream. + * + * @author Gary Russell + * @since 3.2 + * + */ +public final class StreamContainerUtils { + + private StreamContainerUtils() { + } + + /** + * Create a {@link StreamListenerContainer}. + * + * @param consumerDestination the destination. + * @param group the group. + * @param properties the properties. + * @param destination the destination. + * @param extension the properties extension. + * @param applicationContext the application context. + * @return the container. + */ + public static MessageListenerContainer createContainer(ConsumerDestination consumerDestination, String group, + ExtendedConsumerProperties properties, String destination, + RabbitConsumerProperties extension, AbstractApplicationContext applicationContext) { + + StreamListenerContainer container = new StreamListenerContainer(applicationContext.getBean(Environment.class)) { + + @Override + public synchronized void setConsumerCustomizer(ConsumerCustomizer consumerCustomizer) { + super.setConsumerCustomizer((id, builder) -> { + builder.name(consumerDestination.getName() + "." + group); + consumerCustomizer.accept(id, builder); + }); + } + + + }; + container.setBeanName(consumerDestination.getName() + "." + group + ".container"); + container.setMessageConverter(new DefaultStreamMessageConverter()); + return container; + } + + /** + * Configure the channel adapter for streams support. + * @param adapter the adapter. + */ + public static void configureAdapter(AmqpInboundChannelAdapter adapter) { + adapter.setHeaderMapper(new AmqpHeaderMapper() { + + AmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.inboundMapper(); + + @Override + public Map toHeadersFromRequest(MessageProperties source) { + Map headers = this.mapper.toHeadersFromRequest(source); + headers.put("rabbitmq_streamContext", ((StreamMessageProperties) source).getContext()); + return headers; + } + + @Override + public Map toHeadersFromReply(MessageProperties source) { + return null; + } + + @Override + public void fromHeadersToRequest(MessageHeaders headers, MessageProperties target) { + } + + @Override + public void fromHeadersToReply(MessageHeaders headers, MessageProperties target) { + } + + }); + } + +} + +/** + * Temporary work-around for a bug in spring-rabbit-stream 2.4.0-M1. + */ +class DefaultStreamMessageConverter implements StreamMessageConverter { + + private final Supplier builderSupplier; + + private final Charset charset = StandardCharsets.UTF_8; + + /** + * Construct an instance using a {@link WrapperMessageBuilder}. + */ + DefaultStreamMessageConverter() { + this.builderSupplier = () -> new WrapperMessageBuilder(); + } + + /** + * Construct an instance using the provided codec. + * @param codec the codec. + */ + DefaultStreamMessageConverter(@Nullable Codec codec) { + this.builderSupplier = () -> codec.messageBuilder(); + } + + @Override + public Message toMessage(Object object, StreamMessageProperties messageProperties) throws MessageConversionException { + Assert.isInstanceOf(com.rabbitmq.stream.Message.class, object); + com.rabbitmq.stream.Message streamMessage = (com.rabbitmq.stream.Message) object; + toMessageProperties(streamMessage, messageProperties); + return org.springframework.amqp.core.MessageBuilder.withBody(streamMessage.getBodyAsBinary()) + .andProperties(messageProperties) + .build(); + } + + @Override + public com.rabbitmq.stream.Message fromMessage(Message message) throws MessageConversionException { + MessageBuilder builder = this.builderSupplier.get(); + PropertiesBuilder propsBuilder = builder.properties(); + MessageProperties props = message.getMessageProperties(); + Assert.isInstanceOf(StreamMessageProperties.class, props); + StreamMessageProperties mProps = (StreamMessageProperties) props; + JavaUtils.INSTANCE + .acceptIfNotNull(mProps.getMessageId(), propsBuilder::messageId) // TODO different types + .acceptIfNotNull(mProps.getUserId(), usr -> propsBuilder.userId(usr.getBytes(this.charset))) + .acceptIfNotNull(mProps.getTo(), propsBuilder::to) + .acceptIfNotNull(mProps.getSubject(), propsBuilder::subject) + .acceptIfNotNull(mProps.getReplyTo(), propsBuilder::replyTo) + .acceptIfNotNull(mProps.getCorrelationId(), propsBuilder::correlationId) // TODO different types + .acceptIfNotNull(mProps.getContentType(), propsBuilder::contentType) + .acceptIfNotNull(mProps.getContentEncoding(), propsBuilder::contentEncoding) + .acceptIfNotNull(mProps.getExpiration(), exp -> propsBuilder.absoluteExpiryTime(Long.parseLong(exp))) + .acceptIfNotNull(mProps.getCreationTime(), propsBuilder::creationTime) + .acceptIfNotNull(mProps.getGroupId(), propsBuilder::groupId) + .acceptIfNotNull(mProps.getGroupSequence(), propsBuilder::groupSequence) + .acceptIfNotNull(mProps.getReplyToGroupId(), propsBuilder::replyToGroupId); + if (mProps.getHeaders().size() > 0) { + ApplicationPropertiesBuilder appPropsBuilder = builder.applicationProperties(); + mProps.getHeaders().forEach((key, val) -> { + mapProp(key, val, appPropsBuilder); + }); + } + builder.addData(message.getBody()); + return builder.build(); + } + + private void mapProp(String key, Object val, ApplicationPropertiesBuilder builder) { // NOSONAR - complexity + if (val instanceof String) { + builder.entry(key, (String) val); + } + else if (val instanceof Long) { + builder.entry(key, (Long) val); + } + else if (val instanceof Integer) { + builder.entry(key, (Integer) val); + } + else if (val instanceof Short) { + builder.entry(key, (Short) val); + } + else if (val instanceof Byte) { + builder.entry(key, (Byte) val); + } + else if (val instanceof Double) { + builder.entry(key, (Double) val); + } + else if (val instanceof Float) { + builder.entry(key, (Float) val); + } + else if (val instanceof Character) { + builder.entry(key, (Character) val); + } + else if (val instanceof UUID) { + builder.entry(key, (UUID) val); + } + else if (val instanceof byte[]) { + builder.entry(key, (byte[]) val); + } + } + + private void toMessageProperties(com.rabbitmq.stream.Message streamMessage, + StreamMessageProperties mProps) { + + Properties properties = streamMessage.getProperties(); + if (properties != null) { + JavaUtils.INSTANCE + .acceptIfNotNull(properties.getMessageIdAsString(), mProps::setMessageId) + .acceptIfNotNull(properties.getUserId(), usr -> mProps.setUserId(new String(usr, this.charset))) + .acceptIfNotNull(properties.getTo(), mProps::setTo) + .acceptIfNotNull(properties.getSubject(), mProps::setSubject) + .acceptIfNotNull(properties.getReplyTo(), mProps::setReplyTo) + .acceptIfNotNull(properties.getCorrelationIdAsString(), mProps::setCorrelationId) + .acceptIfNotNull(properties.getContentType(), mProps::setContentType) + .acceptIfNotNull(properties.getContentEncoding(), mProps::setContentEncoding) + .acceptIfNotNull(properties.getAbsoluteExpiryTime(), + exp -> mProps.setExpiration(Long.toString(exp))) + .acceptIfNotNull(properties.getCreationTime(), mProps::setCreationTime) + .acceptIfNotNull(properties.getGroupId(), mProps::setGroupId) + .acceptIfNotNull(properties.getGroupSequence(), mProps::setGroupSequence) + .acceptIfNotNull(properties.getReplyToGroupId(), mProps::setReplyToGroupId); + } + Map applicationProperties = streamMessage.getApplicationProperties(); + if (applicationProperties != null) { + mProps.getHeaders().putAll(applicationProperties); + } + } + +} diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java deleted file mode 100644 index a4b412a01..000000000 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java +++ /dev/null @@ -1,169 +0,0 @@ -/* - * Copyright 2021-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. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.binder.rabbit; - -import java.util.Map; -import java.util.function.Function; - -import com.rabbitmq.stream.Environment; - -import org.springframework.amqp.core.MessageProperties; -import org.springframework.amqp.rabbit.listener.MessageListenerContainer; -import org.springframework.amqp.support.converter.MessageConverter; -import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; -import org.springframework.cloud.stream.binder.ExtendedProducerProperties; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties.ProducerType; -import org.springframework.cloud.stream.provisioning.ConsumerDestination; -import org.springframework.cloud.stream.provisioning.ProducerDestination; -import org.springframework.context.support.AbstractApplicationContext; -import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; -import org.springframework.integration.amqp.support.AmqpHeaderMapper; -import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; -import org.springframework.messaging.MessageChannel; -import org.springframework.messaging.MessageHandler; -import org.springframework.messaging.MessageHandlingException; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.support.ErrorMessage; -import org.springframework.rabbit.stream.listener.ConsumerCustomizer; -import org.springframework.rabbit.stream.listener.StreamListenerContainer; -import org.springframework.rabbit.stream.producer.RabbitStreamTemplate; -import org.springframework.rabbit.stream.support.StreamMessageProperties; -import org.springframework.rabbit.stream.support.converter.StreamMessageConverter; - -/** - * Utilities for stream components. Used to prevent a hard runtime dependency on - * spring-rabbit-stream. - * - * @author Gary Russell - * @since 3.2 - * - */ -public final class StreamUtils { - - private StreamUtils() { - } - - /** - * Create a {@link StreamListenerContainer}. - * - * @param consumerDestination the destination. - * @param group the group. - * @param properties the properties. - * @param destination the destination. - * @param extension the properties extension. - * @param applicationContext the application context. - * @return the container. - */ - public static MessageListenerContainer createContainer(ConsumerDestination consumerDestination, String group, - ExtendedConsumerProperties properties, String destination, - RabbitConsumerProperties extension, AbstractApplicationContext applicationContext) { - - StreamListenerContainer container = new StreamListenerContainer(applicationContext.getBean(Environment.class)) { - - @Override - public synchronized void setConsumerCustomizer(ConsumerCustomizer consumerCustomizer) { - super.setConsumerCustomizer((id, builder) -> { - builder.name(consumerDestination.getName() + "." + group); - consumerCustomizer.accept(id, builder); - }); - } - - - }; - container.setBeanName(consumerDestination.getName() + "." + group + ".container"); - String beanName = extension.getStreamStreamMessageConverterBeanName(); - if (beanName != null) { - container.setMessageConverter(applicationContext.getBean(beanName, StreamMessageConverter.class)); - } - return container; - } - - /** - * Configure the channel adapter for streams support. - * @param adapter the adapter. - */ - public static void configureAdapter(AmqpInboundChannelAdapter adapter) { - adapter.setHeaderMapper(new AmqpHeaderMapper() { - - AmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.inboundMapper(); - - @Override - public Map toHeadersFromRequest(MessageProperties source) { - Map headers = this.mapper.toHeadersFromRequest(source); - headers.put("rabbitmq_streamContext", ((StreamMessageProperties) source).getContext()); - return headers; - } - - @Override - public Map toHeadersFromReply(MessageProperties source) { - return null; - } - - @Override - public void fromHeadersToRequest(MessageHeaders headers, MessageProperties target) { - } - - @Override - public void fromHeadersToReply(MessageHeaders headers, MessageProperties target) { - } - - }); - } - - /** - * Create a {@link RabbitStreamMessageHandler}. - * - * @param producerDestination the destination. - * @param producerProperties the properties. - * @param errorChannel the error channel - * @param destination the destination. - * @param extendedProperties the extended properties. - * @param abstractApplicationContext the application context. - * @param headerMapperFunction the header mapper function. - * @return the handler. - */ - public static MessageHandler createStreamMessageHandler(ProducerDestination producerDestination, - ExtendedProducerProperties producerProperties, MessageChannel errorChannel, - String destination, RabbitProducerProperties extendedProperties, - AbstractApplicationContext applicationContext, - Function headerMapperFunction) { - - RabbitStreamTemplate template = new RabbitStreamTemplate(applicationContext.getBean(Environment.class), - producerDestination.getName()); - String beanName = extendedProperties.getStreamMessageConverterBeanName(); - if (beanName != null) { - template.setMessageConverter(applicationContext.getBean(beanName, MessageConverter.class)); - } - beanName = extendedProperties.getStreamStreamMessageConverterBeanName(); - if (beanName != null) { - template.setStreamConverter(applicationContext.getBean(beanName, StreamMessageConverter.class)); - } - RabbitStreamMessageHandler handler = new RabbitStreamMessageHandler(template); - if (errorChannel != null) { - handler.setFailureCallback((msg, ex) -> { - errorChannel.send(new ErrorMessage(new MessageHandlingException(msg, ex))); - }); - } - handler.setHeaderMapper(headerMapperFunction.apply(extendedProperties)); - handler.setSync(ProducerType.STREAM_SYNC.equals(producerProperties.getExtension().getProducerType())); - return handler; - } - -} - diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/AbstractIntegrationTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/AbstractIntegrationTests.java deleted file mode 100644 index ebcb6d4f6..000000000 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/AbstractIntegrationTests.java +++ /dev/null @@ -1,62 +0,0 @@ -/* - * Copyright 2021-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. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.binder.rabbit.stream; - -import java.time.Duration; - -import org.testcontainers.containers.GenericContainer; -import org.testcontainers.utility.DockerImageName; - -/** - * @author Gary Russell - * @since 3.2 - * - */ -public abstract class AbstractIntegrationTests { - - static final GenericContainer RABBITMQ; - - static { - if (System.getProperty("spring.rabbit.use.local.server") == null) { - String image = "pivotalrabbitmq/rabbitmq-stream"; - String cache = System.getenv().get("IMAGE_CACHE"); - if (cache != null) { - image = cache + image; - } - RABBITMQ = new GenericContainer<>(DockerImageName.parse(image)) - .withExposedPorts(5672, 15672, 5552) - .withStartupTimeout(Duration.ofMinutes(2)); - RABBITMQ.start(); - } - else { - RABBITMQ = null; - } - } - - static int amqpPort() { - return RABBITMQ != null ? RABBITMQ.getMappedPort(5672) : 5672; - } - - static int managementPort() { - return RABBITMQ != null ? RABBITMQ.getMappedPort(15672) : 15672; - } - - static int streamPort() { - return RABBITMQ != null ? RABBITMQ.getMappedPort(5552) : 5552; - } - -} diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java index 02cf5a770..ebe7f9a74 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java @@ -19,8 +19,6 @@ package org.springframework.cloud.stream.binder.rabbit.stream; import com.rabbitmq.stream.ConsumerBuilder; import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.OffsetSpecification; -import com.rabbitmq.stream.ProducerBuilder; - import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; @@ -31,22 +29,15 @@ import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.stream.binder.BinderFactory; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; -import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder; -import org.springframework.cloud.stream.binder.rabbit.RabbitStreamMessageHandler; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties.ContainerType; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties.ProducerType; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; -import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; -import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.MessageChannel; -import org.springframework.messaging.MessageHandler; import org.springframework.rabbit.stream.listener.StreamListenerContainer; import static org.assertj.core.api.Assertions.assertThat; @@ -62,7 +53,7 @@ public class RabbitStreamBinderModuleTests { private ConfigurableApplicationContext context; @AfterEach - void tearDown() { + public void tearDown() { if (context != null) { context.close(); context = null; @@ -70,7 +61,7 @@ public class RabbitStreamBinderModuleTests { } @Test - void testStreamContainer() { + public void testStreamContainer() { context = new SpringApplicationBuilder(SimpleProcessor.class) .web(WebApplicationType.NONE) .run("--server.port=0"); @@ -90,36 +81,11 @@ public class RabbitStreamBinderModuleTests { ((StreamListenerContainer) container).stop(); } - @Test - void testStreamHandler() { - context = new SpringApplicationBuilder(SimpleProcessor.class) - .web(WebApplicationType.NONE) - .run("--server.port=0"); - BinderFactory binderFactory = context.getBean(BinderFactory.class); - RabbitMessageChannelBinder rabbitBinder = (RabbitMessageChannelBinder) binderFactory.getBinder(null, - MessageChannel.class); - RabbitProducerProperties rProps = new RabbitProducerProperties(); - rProps.setProducerType(ProducerType.STREAM_SYNC); - ExtendedProducerProperties props = - new ExtendedProducerProperties(rProps); - Binding binding = rabbitBinder.bindProducer("testStream", new DirectChannel(), props); - Object handler = TestUtils.getPropertyValue(binding, "lifecycle"); - assertThat(handler).isInstanceOf(RabbitStreamMessageHandler.class); - } - - @SpringBootApplication(proxyBeanMethods = false) + @SpringBootApplication public static class SimpleProcessor { @Bean - ProducerMessageHandlerCustomizer handlerCustomizer() { - return (hand, dest) -> { - RabbitStreamMessageHandler handler = (RabbitStreamMessageHandler) hand; - handler.setConfirmTimeout(5000); - }; - } - - @Bean - ListenerContainerCustomizer containerCustomizer() { + public ListenerContainerCustomizer containerCustomizer() { return (cont, dest, group) -> { StreamListenerContainer container = (StreamListenerContainer) cont; container.setConsumerCustomizer((name, builder) -> { @@ -129,10 +95,9 @@ public class RabbitStreamBinderModuleTests { } @Bean - Environment env(ConsumerBuilder consumerBuilder, ProducerBuilder producerBuilder) { + Environment env(ConsumerBuilder builder) { Environment env = mock(Environment.class); - given(env.consumerBuilder()).willReturn(consumerBuilder); - given(env.producerBuilder()).willReturn(producerBuilder); + given(env.consumerBuilder()).willReturn(builder); return env; } @@ -141,11 +106,6 @@ public class RabbitStreamBinderModuleTests { return mock(ConsumerBuilder.class); } - @Bean - ProducerBuilder producerBuilder() { - return mock(ProducerBuilder.class); - } - } } diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamMessageHandlerTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamMessageHandlerTests.java deleted file mode 100644 index 9de5c47ec..000000000 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamMessageHandlerTests.java +++ /dev/null @@ -1,114 +0,0 @@ -/* - * Copyright 2021-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. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.binder.rabbit.stream; - -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicReference; - -import com.rabbitmq.stream.Address; -import com.rabbitmq.stream.Consumer; -import com.rabbitmq.stream.Environment; -import com.rabbitmq.stream.OffsetSpecification; -import org.junit.jupiter.api.Test; - -import org.springframework.cloud.stream.binder.rabbit.RabbitStreamMessageHandler; -import org.springframework.integration.support.MessageBuilder; -import org.springframework.rabbit.stream.producer.RabbitStreamTemplate; - -import static org.assertj.core.api.Assertions.assertThat; - -/** - * @author Gary Russell - * @since 3.2 - * - */ -public class RabbitStreamMessageHandlerTests extends AbstractIntegrationTests { - - @Test - void convertAndSend() throws InterruptedException { - Environment env = Environment.builder() - .lazyInitialization(true) - .addressResolver(add -> new Address("localhost", streamPort())) - .build(); - try { - env.deleteStream("stream.stream"); - } - catch (Exception e) { - } - env.streamCreator().stream("stream.stream").create(); - RabbitStreamTemplate streamTemplate = new RabbitStreamTemplate(env, "stream.stream"); - RabbitStreamMessageHandler handler = new RabbitStreamMessageHandler(streamTemplate); - handler.setSync(true); - handler.handleMessage(MessageBuilder.withPayload("foo") - .setHeader("bar", "baz") - .build()); - CountDownLatch latch = new CountDownLatch(1); - AtomicReference received = new AtomicReference<>(); - Consumer consumer = env.consumerBuilder().stream("stream.stream") - .offset(OffsetSpecification.first()) - .messageHandler((context, msg) -> { - received.set(msg); - latch.countDown(); - }) - .build(); - assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); - assertThat(received.get()).isNotNull(); - assertThat(received.get().getBodyAsBinary()).isEqualTo("foo".getBytes()); - assertThat((String) received.get().getApplicationProperties().get("bar")).isEqualTo("baz"); - consumer.close(); - handler.stop(); - } - - @Test - void sendNative() throws InterruptedException { - Environment env = Environment.builder() - .lazyInitialization(true) - .build(); - try { - env.deleteStream("stream.stream"); - } - catch (Exception e) { - } - env.streamCreator().stream("stream.stream").create(); - RabbitStreamTemplate streamTemplate = new RabbitStreamTemplate(env, "stream.stream"); - RabbitStreamMessageHandler handler = new RabbitStreamMessageHandler(streamTemplate); - handler.setSync(true); - handler.handleMessage(MessageBuilder.withPayload(streamTemplate.messageBuilder() - .addData("foo".getBytes()) - .applicationProperties().entry("bar", "baz") - .messageBuilder() - .build()) - .build()); - CountDownLatch latch = new CountDownLatch(1); - AtomicReference received = new AtomicReference<>(); - Consumer consumer = env.consumerBuilder().stream("stream.stream") - .offset(OffsetSpecification.first()) - .messageHandler((context, msg) -> { - received.set(msg); - latch.countDown(); - }) - .build(); - assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); - assertThat(received.get()).isNotNull(); - assertThat(received.get().getBodyAsBinary()).isEqualTo("foo".getBytes()); - assertThat((String) received.get().getApplicationProperties().get("bar")).isEqualTo("baz"); - consumer.close(); - handler.stop(); - } - -}