From 8df13b4709c26af01c5d693cee3b2822ec08027f Mon Sep 17 00:00:00 2001 From: Andy Wilkinson Date: Fri, 30 Jan 2026 11:23:09 +0000 Subject: [PATCH] Polish "Add RabbitMQ Stream service connection from RabbitMQContainer" See gh-42443 --- .../connection/ContainerConnectionSource.java | 7 +- .../pages/testing/testcontainers.adoc | 10 +- ...nectionDetailsFactoryIntegrationTests.java | 1 - ...nectionDetailsFactoryIntegrationTests.java | 32 ++-- ...nectionDetailsFactoryIntegrationTests.java | 167 ++++++++++++++++++ ...nectionDetailsFactoryIntegrationTests.java | 161 +++++++++++++++++ .../RabbitStreamConfiguration.java | 9 +- .../RabbitStreamConnectionDetails.java | 24 +-- ...bbitContainerConnectionDetailsFactory.java | 2 +- ...reamContainerConnectionDetailsFactory.java | 15 +- 10 files changed, 390 insertions(+), 38 deletions(-) create mode 100644 module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SeparateContainersRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java create mode 100644 module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SingleContainerRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java diff --git a/core/spring-boot-testcontainers/src/main/java/org/springframework/boot/testcontainers/service/connection/ContainerConnectionSource.java b/core/spring-boot-testcontainers/src/main/java/org/springframework/boot/testcontainers/service/connection/ContainerConnectionSource.java index f08067721e2..76a23ab82b9 100644 --- a/core/spring-boot-testcontainers/src/main/java/org/springframework/boot/testcontainers/service/connection/ContainerConnectionSource.java +++ b/core/spring-boot-testcontainers/src/main/java/org/springframework/boot/testcontainers/service/connection/ContainerConnectionSource.java @@ -167,7 +167,12 @@ public final class ContainerConnectionSource> implements return this.containerSupplier; } - Set> getConnectionDetailsTypes() { + /** + * Returns the requested connection details types. + * @return the requested connection details types. + * @since 4.1.0 + */ + public Set> getConnectionDetailsTypes() { return this.connectionDetailsTypes; } diff --git a/documentation/spring-boot-docs/src/docs/antora/modules/reference/pages/testing/testcontainers.adoc b/documentation/spring-boot-docs/src/docs/antora/modules/reference/pages/testing/testcontainers.adoc index 4df927c0a0b..4a71a35fc80 100644 --- a/documentation/spring-boot-docs/src/docs/antora/modules/reference/pages/testing/testcontainers.adoc +++ b/documentation/spring-boot-docs/src/docs/antora/modules/reference/pages/testing/testcontainers.adoc @@ -176,6 +176,9 @@ javadoc:org.testcontainers.oracle.OracleContainer[OracleContainer (free)], javad | javadoc:org.springframework.boot.amqp.autoconfigure.RabbitConnectionDetails[] | Containers of type javadoc:{url-testcontainers-rabbitmq-javadoc}/org.testcontainers.rabbitmq.RabbitMQContainer[] +| javadoc:org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails[] +| Containers of type javadoc:{url-testcontainers-rabbitmq-javadoc}/org.testcontainers.rabbitmq.RabbitMQContainer[] when the `@ServiceConnection` `type` attribute includes javadoc:org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails[] + | javadoc:org.springframework.boot.data.redis.autoconfigure.DataRedisConnectionDetails[] | Containers of type javadoc:com.redis.testcontainers.RedisContainer[] or javadoc:com.redis.testcontainers.RedisStackContainer[], or containers named "redis", "redis/redis-stack" or "redis/redis-stack-server" @@ -185,10 +188,13 @@ javadoc:org.testcontainers.oracle.OracleContainer[OracleContainer (free)], javad [TIP] ==== -By default all applicable connection details beans will be created for a given javadoc:org.testcontainers.containers.Container[]. +By default, with the exception of javadoc:org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails[], all applicable connection details beans will be created for a given javadoc:org.testcontainers.containers.Container[]. For example, a javadoc:{url-testcontainers-postgresql-javadoc}/org.testcontainers.postgresql.PostgreSQLContainer[] will create both javadoc:org.springframework.boot.jdbc.autoconfigure.JdbcConnectionDetails[] and javadoc:org.springframework.boot.r2dbc.autoconfigure.R2dbcConnectionDetails[]. If you want to create only a subset of the applicable types, you can use the `type` attribute of javadoc:org.springframework.boot.testcontainers.service.connection.ServiceConnection[format=annotation]. + +To create a javadoc:org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails[] bean from a javadoc:{url-testcontainers-rabbitmq-javadoc}/org.testcontainers.rabbitmq.RabbitMQContainer[], you must opt in using the `type` attribute of javadoc:org.springframework.boot.testcontainers.service.connection.ServiceConnection[format=annotation]. +The container must also expose port 5552, the RabbitMQ streams port. ==== By default `Container.getDockerImageName().getRepository()` is used to obtain the name used to find connection details. @@ -229,7 +235,7 @@ The SSL annotations are supported for the following service connections: * Elasticsearch * Kafka * MongoDB -* RabbitMQ +* RabbitMQ (excluding streams) * Redis The `ElasticsearchContainer` additionally supports automatic detection of server side SSL. diff --git a/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactoryIntegrationTests.java b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactoryIntegrationTests.java index 898f9fad3dc..ea80da8d9bf 100644 --- a/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactoryIntegrationTests.java +++ b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactoryIntegrationTests.java @@ -71,7 +71,6 @@ class RabbitContainerConnectionDetailsFactoryIntegrationTests { this.rabbitTemplate.convertAndSend("test", "message"); Awaitility.waitAtMost(Duration.ofMinutes(4)) .untilAsserted(() -> assertThat(this.listener.messages).containsExactly("message")); - } @Configuration(proxyBeanMethods = false) diff --git a/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java index 53d1117e399..1e65d213f2a 100644 --- a/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java +++ b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2024 the original author or authors. + * Copyright 2012-present 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. @@ -24,16 +24,19 @@ import com.rabbitmq.stream.Address; import com.rabbitmq.stream.Environment; import org.awaitility.Awaitility; import org.junit.jupiter.api.Test; -import org.testcontainers.containers.RabbitMQContainer; import org.testcontainers.images.builder.Transferable; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.rabbitmq.RabbitMQContainer; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.amqp.autoconfigure.EnvironmentBuilderCustomizer; import org.springframework.boot.amqp.autoconfigure.RabbitAutoConfiguration; +import org.springframework.boot.amqp.autoconfigure.RabbitConnectionDetails; import org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails; +import org.springframework.boot.amqp.testcontainers.RabbitContainerConnectionDetailsFactory.RabbitMqContainerConnectionDetails; +import org.springframework.boot.amqp.testcontainers.RabbitStreamContainerConnectionDetailsFactory.RabbitMqStreamContainerConnectionDetails; import org.springframework.boot.autoconfigure.ImportAutoConfiguration; import org.springframework.boot.testcontainers.service.connection.ServiceConnection; import org.springframework.boot.testsupport.container.TestImage; @@ -47,9 +50,11 @@ import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; import static org.assertj.core.api.Assertions.assertThat; /** - * Tests for {@link RabbitStreamContainerConnectionDetailsFactory}. + * Tests for {@link RabbitStreamContainerConnectionDetailsFactory} with a single container + * that's only used for streams. * * @author Eddú Meléndez + * @author Andy Wilkinson */ @SpringJUnitConfig @TestPropertySource( @@ -60,7 +65,7 @@ class RabbitStreamContainerConnectionDetailsFactoryIntegrationTests { private static final int RABBITMQ_STREAMS_PORT = 5552; @Container - @ServiceConnection + @ServiceConnection(type = RabbitStreamConnectionDetails.class) static final RabbitMQContainer rabbit = getRabbitMqStreamContainer(); private static RabbitMQContainer getRabbitMqStreamContainer() { @@ -72,7 +77,10 @@ class RabbitStreamContainerConnectionDetailsFactoryIntegrationTests { } @Autowired(required = false) - private RabbitStreamConnectionDetails connectionDetails; + private RabbitConnectionDetails connectionDetails; + + @Autowired(required = false) + private RabbitStreamConnectionDetails streamConnectionDetails; @Autowired private RabbitStreamTemplate rabbitStreamTemplate; @@ -82,11 +90,11 @@ class RabbitStreamContainerConnectionDetailsFactoryIntegrationTests { @Test void connectionCanBeMadeToRabbitContainer() { - assertThat(this.connectionDetails).isNotNull(); + assertThat(this.connectionDetails).isNotInstanceOf(RabbitMqContainerConnectionDetails.class); + assertThat(this.streamConnectionDetails).isInstanceOf(RabbitMqStreamContainerConnectionDetails.class); this.rabbitStreamTemplate.convertAndSend("message"); Awaitility.waitAtMost(Duration.ofMinutes(4)) .untilAsserted(() -> assertThat(this.listener.messages).containsExactly("message")); - } @Configuration(proxyBeanMethods = false) @@ -95,17 +103,13 @@ class RabbitStreamContainerConnectionDetailsFactoryIntegrationTests { @Bean StreamAdmin streamAdmin(Environment env) { - return new StreamAdmin(env, sc -> { - sc.stream("stream.queue1").create(); - }); + return new StreamAdmin(env, (sc) -> sc.stream("stream.queue1").create()); } @Bean EnvironmentBuilderCustomizer environmentBuilderCustomizer() { - return env -> { - Address entrypoint = new Address(rabbit.getHost(), rabbit.getMappedPort(RABBITMQ_STREAMS_PORT)); - env.addressResolver(address -> entrypoint); - }; + return (env) -> env.addressResolver( + (address) -> new Address(rabbit.getHost(), rabbit.getMappedPort(RABBITMQ_STREAMS_PORT))); } @Bean diff --git a/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SeparateContainersRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SeparateContainersRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java new file mode 100644 index 00000000000..778f6536bd5 --- /dev/null +++ b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SeparateContainersRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java @@ -0,0 +1,167 @@ +/* + * Copyright 2012-present 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.boot.amqp.testcontainers; + +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; + +import com.rabbitmq.stream.Address; +import com.rabbitmq.stream.Environment; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.Test; +import org.testcontainers.images.builder.Transferable; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.rabbitmq.RabbitMQContainer; + +import org.springframework.amqp.core.AmqpAdmin; +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.amqp.autoconfigure.EnvironmentBuilderCustomizer; +import org.springframework.boot.amqp.autoconfigure.RabbitAutoConfiguration; +import org.springframework.boot.amqp.autoconfigure.RabbitConnectionDetails; +import org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails; +import org.springframework.boot.amqp.testcontainers.RabbitContainerConnectionDetailsFactory.RabbitMqContainerConnectionDetails; +import org.springframework.boot.amqp.testcontainers.RabbitStreamContainerConnectionDetailsFactory.RabbitMqStreamContainerConnectionDetails; +import org.springframework.boot.autoconfigure.ImportAutoConfiguration; +import org.springframework.boot.testcontainers.service.connection.ServiceConnection; +import org.springframework.boot.testsupport.container.TestImage; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.rabbit.stream.producer.RabbitStreamTemplate; +import org.springframework.rabbit.stream.support.StreamAdmin; +import org.springframework.test.context.TestPropertySource; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests for {@link RabbitStreamContainerConnectionDetailsFactory} with two containers, + * one for streams and one for standard messaging. + * + * @author Eddú Meléndez + * @author Andy Wilkinson + */ +@SpringJUnitConfig +@TestPropertySource( + properties = { "spring.rabbitmq.stream.name=stream.queue1", "spring.rabbitmq.listener.type=stream" }) +@Testcontainers(disabledWithoutDocker = true) +class SeparateContainersRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests { + + private static final int RABBITMQ_STREAMS_PORT = 5552; + + @Container + @ServiceConnection + static final RabbitMQContainer rabbit = TestImage.container(RabbitMQContainer.class); + + @Container + @ServiceConnection(type = RabbitStreamConnectionDetails.class) + static final RabbitMQContainer rabbitStream = getRabbitMqStreamContainer(); + + private static RabbitMQContainer getRabbitMqStreamContainer() { + RabbitMQContainer container = TestImage.container(RabbitMQContainer.class); + container.addExposedPorts(RABBITMQ_STREAMS_PORT); + String enabledPlugins = "[rabbitmq_stream,rabbitmq_prometheus]."; + container.withCopyToContainer(Transferable.of(enabledPlugins), "/etc/rabbitmq/enabled_plugins"); + return container; + } + + @Autowired(required = false) + private RabbitConnectionDetails connectionDetails; + + @Autowired(required = false) + private RabbitStreamConnectionDetails streamConnectionDetails; + + @Autowired + private RabbitTemplate rabbitTemplate; + + @Autowired + private RabbitStreamTemplate rabbitStreamTemplate; + + @Autowired + private StreamTestListener streamListener; + + @Test + void rabbitConnectionDetailsAreSourcedFromContainer() { + assertThat(this.connectionDetails).isInstanceOf(RabbitMqContainerConnectionDetails.class); + assertThat(this.connectionDetails.getFirstAddress().port()).isEqualTo(rabbit.getAmqpPort()); + } + + @Test + void rabbitStreamConnectionDetailsAreSourcedFromContainer() { + assertThat(this.streamConnectionDetails).isInstanceOf(RabbitMqStreamContainerConnectionDetails.class); + assertThat(this.streamConnectionDetails.getPort()).isEqualTo(rabbitStream.getMappedPort(RABBITMQ_STREAMS_PORT)); + } + + @Test + void connectionCanBeMadeToRabbitContainer() { + this.rabbitTemplate.convertAndSend("test", "message"); + Awaitility.waitAtMost(Duration.ofMinutes(4)) + .untilAsserted(() -> assertThat(this.rabbitTemplate.receive("test")) + .extracting((message) -> new String(message.getBody(), StandardCharsets.UTF_8)) + .isEqualTo("message")); + } + + @Test + void streamConnectionCanBeMadeToRabbitContainer() { + this.rabbitStreamTemplate.convertAndSend("message"); + Awaitility.waitAtMost(Duration.ofMinutes(4)) + .untilAsserted(() -> assertThat(this.streamListener.messages).containsExactly("message")); + } + + @Configuration(proxyBeanMethods = false) + @ImportAutoConfiguration(RabbitAutoConfiguration.class) + static class TestConfiguration { + + TestConfiguration(AmqpAdmin amqpAdmin) { + amqpAdmin.declareQueue(new Queue("test")); + } + + @Bean + StreamAdmin streamAdmin(Environment env) { + return new StreamAdmin(env, (sc) -> sc.stream("stream.queue1").create()); + } + + @Bean + EnvironmentBuilderCustomizer environmentBuilderCustomizer() { + return (env) -> env.addressResolver((address) -> new Address(rabbitStream.getHost(), + rabbitStream.getMappedPort(RABBITMQ_STREAMS_PORT))); + } + + @Bean + StreamTestListener streamTestListener() { + return new StreamTestListener(); + } + + } + + static class StreamTestListener { + + private final List messages = new ArrayList<>(); + + @RabbitListener(queues = "stream.queue1") + void processMessage(String message) { + this.messages.add(message); + } + + } + +} diff --git a/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SingleContainerRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SingleContainerRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java new file mode 100644 index 00000000000..6dc8b12291c --- /dev/null +++ b/module/spring-boot-amqp/src/dockerTest/java/org/springframework/boot/amqp/testcontainers/SingleContainerRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests.java @@ -0,0 +1,161 @@ +/* + * Copyright 2012-present 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.boot.amqp.testcontainers; + +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; + +import com.rabbitmq.stream.Address; +import com.rabbitmq.stream.Environment; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.Test; +import org.testcontainers.images.builder.Transferable; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.rabbitmq.RabbitMQContainer; + +import org.springframework.amqp.core.AmqpAdmin; +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.amqp.autoconfigure.EnvironmentBuilderCustomizer; +import org.springframework.boot.amqp.autoconfigure.RabbitAutoConfiguration; +import org.springframework.boot.amqp.autoconfigure.RabbitConnectionDetails; +import org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails; +import org.springframework.boot.amqp.testcontainers.RabbitContainerConnectionDetailsFactory.RabbitMqContainerConnectionDetails; +import org.springframework.boot.amqp.testcontainers.RabbitStreamContainerConnectionDetailsFactory.RabbitMqStreamContainerConnectionDetails; +import org.springframework.boot.autoconfigure.ImportAutoConfiguration; +import org.springframework.boot.testcontainers.service.connection.ServiceConnection; +import org.springframework.boot.testsupport.container.TestImage; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.rabbit.stream.producer.RabbitStreamTemplate; +import org.springframework.rabbit.stream.support.StreamAdmin; +import org.springframework.test.context.TestPropertySource; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests for {@link RabbitStreamContainerConnectionDetailsFactory} with a single container + * used for both streams and standard messaging. + * + * @author Eddú Meléndez + * @author Andy Wilkinson + */ +@SpringJUnitConfig +@TestPropertySource( + properties = { "spring.rabbitmq.stream.name=stream.queue1", "spring.rabbitmq.listener.type=stream" }) +@Testcontainers(disabledWithoutDocker = true) +class SingleContainerRabbitAndRabbitStreamContainerConnectionDetailsFactoryIntegrationTests { + + private static final int RABBITMQ_STREAMS_PORT = 5552; + + @Container + @ServiceConnection(type = { RabbitConnectionDetails.class, RabbitStreamConnectionDetails.class }) + static final RabbitMQContainer rabbit = getRabbitMqStreamContainer(); + + private static RabbitMQContainer getRabbitMqStreamContainer() { + RabbitMQContainer container = TestImage.container(RabbitMQContainer.class); + container.addExposedPorts(RABBITMQ_STREAMS_PORT); + String enabledPlugins = "[rabbitmq_stream,rabbitmq_prometheus]."; + container.withCopyToContainer(Transferable.of(enabledPlugins), "/etc/rabbitmq/enabled_plugins"); + return container; + } + + @Autowired(required = false) + private RabbitConnectionDetails connectionDetails; + + @Autowired(required = false) + private RabbitStreamConnectionDetails streamConnectionDetails; + + @Autowired + private RabbitStreamTemplate rabbitStreamTemplate; + + @Autowired + private RabbitTemplate rabbitTemplate; + + @Autowired + private StreamTestListener streamListener; + + @Test + void rabbitConnectionDetailsAreSourcedFromContainer() { + assertThat(this.connectionDetails).isInstanceOf(RabbitMqContainerConnectionDetails.class); + } + + @Test + void rabbitStreamConnectionDetailsAreSourcedFromContainer() { + assertThat(this.streamConnectionDetails).isInstanceOf(RabbitMqStreamContainerConnectionDetails.class); + } + + @Test + void connectionCanBeMadeToRabbitContainer() { + this.rabbitTemplate.convertAndSend("test", "message"); + Awaitility.waitAtMost(Duration.ofMinutes(4)) + .untilAsserted(() -> assertThat(this.rabbitTemplate.receive("test")) + .extracting((message) -> new String(message.getBody(), StandardCharsets.UTF_8)) + .isEqualTo("message")); + } + + @Test + void streamConnectionCanBeMadeToRabbitContainer() { + this.rabbitStreamTemplate.convertAndSend("message"); + Awaitility.waitAtMost(Duration.ofMinutes(4)) + .untilAsserted(() -> assertThat(this.streamListener.messages).containsExactly("message")); + } + + @Configuration(proxyBeanMethods = false) + @ImportAutoConfiguration(RabbitAutoConfiguration.class) + static class TestConfiguration { + + TestConfiguration(AmqpAdmin amqpAdmin) { + amqpAdmin.declareQueue(new Queue("test")); + } + + @Bean + StreamAdmin streamAdmin(Environment env) { + return new StreamAdmin(env, (sc) -> sc.stream("stream.queue1").create()); + } + + @Bean + EnvironmentBuilderCustomizer environmentBuilderCustomizer() { + return (env) -> env.addressResolver( + (address) -> new Address(rabbit.getHost(), rabbit.getMappedPort(RABBITMQ_STREAMS_PORT))); + } + + @Bean + StreamTestListener streamTestListener() { + return new StreamTestListener(); + } + + } + + static class StreamTestListener { + + private final List messages = new ArrayList<>(); + + @RabbitListener(queues = "stream.queue1") + void processMessage(String message) { + this.messages.add(message); + } + + } + +} diff --git a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConfiguration.java b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConfiguration.java index b06fa1909c7..0b396dd01f7 100644 --- a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConfiguration.java +++ b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConfiguration.java @@ -18,6 +18,7 @@ package org.springframework.boot.amqp.autoconfigure; import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.EnvironmentBuilder; +import org.jspecify.annotations.Nullable; import org.springframework.amqp.rabbit.config.ContainerCustomizer; import org.springframework.amqp.support.converter.MessageConverter; @@ -52,7 +53,7 @@ import org.springframework.util.Assert; class RabbitStreamConfiguration { @Bean - @ConditionalOnMissingBean(RabbitStreamConnectionDetails.class) + @ConditionalOnMissingBean RabbitStreamConnectionDetails rabbitStreamConnectionDetails(RabbitProperties rabbitProperties) { return new PropertiesRabbitStreamConnectionDetails(rabbitProperties.getStream()); } @@ -152,17 +153,17 @@ class RabbitStreamConfiguration { } @Override - public String getVirtualHost() { + public @Nullable String getVirtualHost() { return this.streamProperties.getVirtualHost(); } @Override - public String getUsername() { + public @Nullable String getUsername() { return this.streamProperties.getUsername(); } @Override - public String getPassword() { + public @Nullable String getPassword() { return this.streamProperties.getPassword(); } diff --git a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConnectionDetails.java b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConnectionDetails.java index de630e6f9d9..8b08349167d 100644 --- a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConnectionDetails.java +++ b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/autoconfigure/RabbitStreamConnectionDetails.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2024 the original author or authors. + * Copyright 2012-present 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. @@ -16,6 +16,8 @@ package org.springframework.boot.amqp.autoconfigure; +import org.jspecify.annotations.Nullable; + import org.springframework.boot.autoconfigure.service.connection.ConnectionDetails; /** @@ -27,30 +29,30 @@ import org.springframework.boot.autoconfigure.service.connection.ConnectionDetai public interface RabbitStreamConnectionDetails extends ConnectionDetails { /** - * Rabbit server host. - * @return the rabbit server host + * Rabbit Stream server host. + * @return the Rabbit Stream server host */ String getHost(); /** * Rabbit Stream server port. - * @return the rabbit stream server port + * @return the Rabbit Stream server port */ int getPort(); /** - * Login user to authenticate to the broker. - * @return the login user to authenticate to the broker or {@code null} + * Username for authentication. + * @return the username for authentication or {@code null} */ - default String getUsername() { + default @Nullable String getUsername() { return null; } /** - * Login to authenticate against the broker. - * @return the login to authenticate against the broker or {@code null} + * Password for authentication. + * @return the password for authentication or {@code null} */ - default String getPassword() { + default @Nullable String getPassword() { return null; } @@ -58,7 +60,7 @@ public interface RabbitStreamConnectionDetails extends ConnectionDetails { * Virtual host to use when connecting to the broker. * @return the virtual host to use when connecting to the broker or {@code null} */ - default String getVirtualHost() { + default @Nullable String getVirtualHost() { return null; } diff --git a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactory.java b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactory.java index 8049dd52218..fd47f98fd47 100644 --- a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactory.java +++ b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitContainerConnectionDetailsFactory.java @@ -49,7 +49,7 @@ class RabbitContainerConnectionDetailsFactory /** * {@link RabbitConnectionDetails} backed by a {@link ContainerConnectionSource}. */ - private static final class RabbitMqContainerConnectionDetails extends ContainerConnectionDetails + static final class RabbitMqContainerConnectionDetails extends ContainerConnectionDetails implements RabbitConnectionDetails { private RabbitMqContainerConnectionDetails(ContainerConnectionSource source) { diff --git a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactory.java b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactory.java index fb364a90703..8d2bf2bc9dc 100644 --- a/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactory.java +++ b/module/spring-boot-amqp/src/main/java/org/springframework/boot/amqp/testcontainers/RabbitStreamContainerConnectionDetailsFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2024 the original author or authors. + * Copyright 2012-present 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. @@ -16,7 +16,7 @@ package org.springframework.boot.amqp.testcontainers; -import org.testcontainers.containers.RabbitMQContainer; +import org.testcontainers.rabbitmq.RabbitMQContainer; import org.springframework.boot.amqp.autoconfigure.RabbitStreamConnectionDetails; import org.springframework.boot.testcontainers.service.connection.ContainerConnectionDetailsFactory; @@ -37,6 +37,13 @@ class RabbitStreamContainerConnectionDetailsFactory super(ANY_CONNECTION_NAME, "org.springframework.rabbit.stream.producer.RabbitStreamTemplate"); } + @Override + protected boolean sourceAccepts(ContainerConnectionSource source, Class requiredContainerType, + Class requiredConnectionDetailsType) { + return source.getConnectionDetailsTypes().contains(requiredConnectionDetailsType) + && super.sourceAccepts(source, requiredContainerType, requiredConnectionDetailsType); + } + @Override protected RabbitStreamConnectionDetails getContainerConnectionDetails( ContainerConnectionSource source) { @@ -47,8 +54,8 @@ class RabbitStreamContainerConnectionDetailsFactory * {@link RabbitStreamConnectionDetails} backed by a * {@link ContainerConnectionSource}. */ - private static final class RabbitMqStreamContainerConnectionDetails - extends ContainerConnectionDetails implements RabbitStreamConnectionDetails { + static final class RabbitMqStreamContainerConnectionDetails extends ContainerConnectionDetails + implements RabbitStreamConnectionDetails { private RabbitMqStreamContainerConnectionDetails(ContainerConnectionSource source) { super(source);