From 920b8ae744aa9fb34fccaa124613ecaa36b03f7a Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 7 Jul 2021 15:57:44 -0400 Subject: [PATCH] ZeroMQ test: Add sleep between SUB & PUB --- .../zeromq/outbound/ZeroMqMessageHandlerTests.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/spring-integration-zeromq/src/test/java/org/springframework/integration/zeromq/outbound/ZeroMqMessageHandlerTests.java b/spring-integration-zeromq/src/test/java/org/springframework/integration/zeromq/outbound/ZeroMqMessageHandlerTests.java index ecb1a8b546..8935e124b5 100644 --- a/spring-integration-zeromq/src/test/java/org/springframework/integration/zeromq/outbound/ZeroMqMessageHandlerTests.java +++ b/spring-integration-zeromq/src/test/java/org/springframework/integration/zeromq/outbound/ZeroMqMessageHandlerTests.java @@ -86,7 +86,7 @@ public class ZeroMqMessageHandlerTests { } @Test - void testMessageHandlerForPubSub() { + void testMessageHandlerForPubSub() throws InterruptedException { ZMQ.Socket subSocket = CONTEXT.createSocket(SocketType.SUB); subSocket.setReceiveTimeOut(20_000); int port = subSocket.bindToRandomPort("tcp://*"); @@ -98,12 +98,15 @@ public class ZeroMqMessageHandlerTests { new FunctionExpression>((message) -> message.getHeaders().get("topic"))); messageHandler.setMessageMapper(new EmbeddedJsonHeadersMessageMapper()); messageHandler.afterPropertiesSet(); + subSocket.subscribe("test"); + + // Give it some time to connect and subscribe + Thread.sleep(2000); Message testMessage = MessageBuilder.withPayload("test").setHeader("topic", "testTopic").build(); await().atMost(Duration.ofSeconds(20)).pollDelay(Duration.ofMillis(100)) .untilAsserted(() -> { - subSocket.subscribe("test"); messageHandler.handleMessage(testMessage).subscribe(); ZMsg msg = ZMsg.recvMsg(subSocket); assertThat(msg).isNotNull();