diff options
| author | Arnaud Cogoluègnes <acogoluegnes@gmail.com> | 2020-06-11 18:12:31 +0200 |
|---|---|---|
| committer | Arnaud Cogoluègnes <acogoluegnes@gmail.com> | 2020-06-11 18:12:31 +0200 |
| commit | 50581b22554ec9a8cddcdddc6a27b2a704d541f7 (patch) | |
| tree | c1be3fb1f0e3a93e5d3a9988826f834c6614774d /deps | |
| parent | 157808ca8fa08346b02c5ce9747880009ae5442c (diff) | |
| download | rabbitmq-server-git-50581b22554ec9a8cddcdddc6a27b2a704d541f7.tar.gz | |
Fix cluster tests
Diffstat (limited to 'deps')
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; } |
