summaryrefslogtreecommitdiff
path: root/deps
diff options
context:
space:
mode:
authorArnaud Cogoluègnes <acogoluegnes@gmail.com>2020-06-11 18:12:31 +0200
committerArnaud Cogoluègnes <acogoluegnes@gmail.com>2020-06-11 18:12:31 +0200
commit50581b22554ec9a8cddcdddc6a27b2a704d541f7 (patch)
treec1be3fb1f0e3a93e5d3a9988826f834c6614774d /deps
parent157808ca8fa08346b02c5ce9747880009ae5442c (diff)
downloadrabbitmq-server-git-50581b22554ec9a8cddcdddc6a27b2a704d541f7.tar.gz
Fix cluster tests
Diffstat (limited to 'deps')
-rw-r--r--deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/StreamTest.java7
-rw-r--r--deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/TestUtils.java3
2 files changed, 5 insertions, 5 deletions
diff --git a/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/StreamTest.java b/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/StreamTest.java
index 954d74349b..eb23682438 100644
--- a/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/StreamTest.java
+++ b/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/StreamTest.java
@@ -61,15 +61,15 @@ public class StreamTest {
@MethodSource
void shouldBePossibleToPublishFromAnyNodeAndConsumeFromAnyMember(Function<Client.StreamMetadata, Client.Broker> publisherBroker,
Function<Client.StreamMetadata, Client.Broker> consumerBroker) throws Exception {
+
int messageCount = 10_000;
- Client client = cf.get();
+ Client client = cf.get(new Client.ClientParameters().port(TestUtils.streamPort()));
Map<String, Client.StreamMetadata> metadata = client.metadata(stream);
assertThat(metadata).hasSize(1).containsKey(stream);
Client.StreamMetadata streamMetadata = metadata.get(stream);
CountDownLatch publishingLatch = new CountDownLatch(messageCount);
Client publisher = cf.get(new Client.ClientParameters()
- .host(publisherBroker.apply(streamMetadata).getHost())
.port(publisherBroker.apply(streamMetadata).getPort())
.confirmListener(publishingId -> publishingLatch.countDown()));
@@ -80,7 +80,6 @@ public class StreamTest {
CountDownLatch consumingLatch = new CountDownLatch(messageCount);
Set<String> bodies = ConcurrentHashMap.newKeySet(messageCount);
Client consumer = cf.get(new Client.ClientParameters()
- .host(consumerBroker.apply(streamMetadata).getHost())
.port(consumerBroker.apply(streamMetadata).getPort())
.chunkListener((client1, subscriptionId, offset, messageCount1, dataSize) -> client1.credit(subscriptionId, 10))
.messageListener((subscriptionId, offset, message) -> {
@@ -98,7 +97,7 @@ public class StreamTest {
@Test
void metadataOnClusterShouldReturnLeaderAndReplicas() {
- Client client = cf.get();
+ Client client = cf.get(new Client.ClientParameters().port(TestUtils.streamPort()));
Map<String, Client.StreamMetadata> metadata = client.metadata(stream);
assertThat(metadata).hasSize(1).containsKey(stream);
Client.StreamMetadata streamMetadata = metadata.get(stream);
diff --git a/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/TestUtils.java b/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/TestUtils.java
index 550a426f67..2419cb7c61 100644
--- a/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/TestUtils.java
+++ b/deps/rabbitmq_stream/test/rabbit_stream_SUITE_data/src/test/java/com/rabbitmq/stream/TestUtils.java
@@ -135,7 +135,8 @@ public class TestUtils {
}
public Client get(Client.ClientParameters parameters) {
- Client client = new Client(parameters.eventLoopGroup(eventLoopGroup).port(streamPort()));
+ // don't set the port, it would override the caller's port setting
+ Client client = new Client(parameters.eventLoopGroup(eventLoopGroup));
clients.add(client);
return client;
}