mirror of
https://github.com/spring-projects/spring-boot.git
synced 2026-10-09 19:09:03 +00:00
Polish "Add RabbitMQ Stream service connection from RabbitMQContainer"
See gh-42443
This commit is contained in:
+6
-1
@@ -167,7 +167,12 @@ public final class ContainerConnectionSource<C extends Container<?>> implements
|
||||
return this.containerSupplier;
|
||||
}
|
||||
|
||||
Set<Class<?>> getConnectionDetailsTypes() {
|
||||
/**
|
||||
* Returns the requested connection details types.
|
||||
* @return the requested connection details types.
|
||||
* @since 4.1.0
|
||||
*/
|
||||
public Set<Class<?>> getConnectionDetailsTypes() {
|
||||
return this.connectionDetailsTypes;
|
||||
}
|
||||
|
||||
|
||||
+8
-2
@@ -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.
|
||||
|
||||
-1
@@ -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)
|
||||
|
||||
+18
-14
@@ -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
|
||||
|
||||
+167
@@ -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<String> messages = new ArrayList<>();
|
||||
|
||||
@RabbitListener(queues = "stream.queue1")
|
||||
void processMessage(String message) {
|
||||
this.messages.add(message);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
+161
@@ -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<String> messages = new ArrayList<>();
|
||||
|
||||
@RabbitListener(queues = "stream.queue1")
|
||||
void processMessage(String message) {
|
||||
this.messages.add(message);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
+5
-4
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
+13
-11
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -49,7 +49,7 @@ class RabbitContainerConnectionDetailsFactory
|
||||
/**
|
||||
* {@link RabbitConnectionDetails} backed by a {@link ContainerConnectionSource}.
|
||||
*/
|
||||
private static final class RabbitMqContainerConnectionDetails extends ContainerConnectionDetails<RabbitMQContainer>
|
||||
static final class RabbitMqContainerConnectionDetails extends ContainerConnectionDetails<RabbitMQContainer>
|
||||
implements RabbitConnectionDetails {
|
||||
|
||||
private RabbitMqContainerConnectionDetails(ContainerConnectionSource<RabbitMQContainer> source) {
|
||||
|
||||
+11
-4
@@ -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<RabbitMQContainer> source, Class<?> requiredContainerType,
|
||||
Class<?> requiredConnectionDetailsType) {
|
||||
return source.getConnectionDetailsTypes().contains(requiredConnectionDetailsType)
|
||||
&& super.sourceAccepts(source, requiredContainerType, requiredConnectionDetailsType);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected RabbitStreamConnectionDetails getContainerConnectionDetails(
|
||||
ContainerConnectionSource<RabbitMQContainer> source) {
|
||||
@@ -47,8 +54,8 @@ class RabbitStreamContainerConnectionDetailsFactory
|
||||
* {@link RabbitStreamConnectionDetails} backed by a
|
||||
* {@link ContainerConnectionSource}.
|
||||
*/
|
||||
private static final class RabbitMqStreamContainerConnectionDetails
|
||||
extends ContainerConnectionDetails<RabbitMQContainer> implements RabbitStreamConnectionDetails {
|
||||
static final class RabbitMqStreamContainerConnectionDetails extends ContainerConnectionDetails<RabbitMQContainer>
|
||||
implements RabbitStreamConnectionDetails {
|
||||
|
||||
private RabbitMqStreamContainerConnectionDetails(ContainerConnectionSource<RabbitMQContainer> source) {
|
||||
super(source);
|
||||
|
||||
Reference in New Issue
Block a user