From 96adc08ad867bdd3a267a7a004e9998d6167d11d Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 6 Nov 2023 15:14:52 -0500 Subject: [PATCH] Adapt AMQP tests to latest Rabbit Streams Client Related to https://github.com/spring-projects/spring-amqp/issues/2522 No need to use an `AddressResolver` with the latest RabbitMQ Streams Client library The current configuration is reflecting whatever Spring Boot auto-configuration experience would expect from us --- .../integration/amqp/dsl/RabbitStreamTests.java | 3 +-- .../amqp/outbound/RabbitStreamMessageHandlerTests.java | 5 ++--- .../integration/amqp/support/RabbitTestContainer.java | 1 - 3 files changed, 3 insertions(+), 6 deletions(-) diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/RabbitStreamTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/RabbitStreamTests.java index 9674d8e9d4..d54521493c 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/RabbitStreamTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/RabbitStreamTests.java @@ -16,7 +16,6 @@ package org.springframework.integration.amqp.dsl; -import com.rabbitmq.stream.Address; import com.rabbitmq.stream.Environment; import org.junit.jupiter.api.Test; @@ -103,7 +102,7 @@ public class RabbitStreamTests implements RabbitTestContainer { @Bean Environment rabbitStreamEnvironment() { return Environment.builder() - .addressResolver(add -> new Address("localhost", RabbitTestContainer.streamPort())) + .port(RabbitTestContainer.streamPort()) .build(); } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/RabbitStreamMessageHandlerTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/RabbitStreamMessageHandlerTests.java index 3580fbae9a..71e607db81 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/RabbitStreamMessageHandlerTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/RabbitStreamMessageHandlerTests.java @@ -20,7 +20,6 @@ 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; @@ -46,7 +45,7 @@ public class RabbitStreamMessageHandlerTests implements RabbitTestContainer { void convertAndSend() throws InterruptedException { Environment env = Environment.builder() .lazyInitialization(true) - .addressResolver(add -> new Address("localhost", RabbitTestContainer.streamPort())) + .port(RabbitTestContainer.streamPort()) .build(); try { env.deleteStream("stream.stream"); @@ -83,7 +82,7 @@ public class RabbitStreamMessageHandlerTests implements RabbitTestContainer { @Test void sendNative() throws InterruptedException { Environment env = Environment.builder() - .addressResolver(add -> new Address("localhost", RabbitTestContainer.streamPort())) + .port(RabbitTestContainer.streamPort()) .lazyInitialization(true) .build(); try { diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/RabbitTestContainer.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/RabbitTestContainer.java index cf045cc03b..1ea4891df5 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/RabbitTestContainer.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/RabbitTestContainer.java @@ -35,7 +35,6 @@ public interface RabbitTestContainer { RabbitMQContainer RABBITMQ = new RabbitMQContainer("rabbitmq:management") .withExposedPorts(5672, 15672, 5552) - .withEnv("RABBITMQ_SERVER_ADDITIONAL_ERL_ARGS", "-rabbitmq_stream advertised_host localhost") .withStartupTimeout(Duration.ofMinutes(2)); @BeforeAll