From cd7465a370a950f9d8daad20901ffc8449ace4ab Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 22 Jun 2021 11:36:27 -0400 Subject: [PATCH] Rework ZeroMQMH test for Awaitility Turns out PUB socket doesn't care if there are subscribers to it or not. The sent message may be just lost in between. * Resend message in the test until it is received by subscriber * Use `await().untilAsserted()` to iterate the logic at most 10 seconds --- .../outbound/ZeroMqMessageHandlerTests.java | 27 +++++++------------ 1 file changed, 9 insertions(+), 18 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 296fdb496d..ceaba6bec5 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 @@ -100,27 +100,18 @@ public class ZeroMqMessageHandlerTests { messageHandler.setMessageMapper(new EmbeddedJsonHeadersMessageMapper()); messageHandler.afterPropertiesSet(); - ZMQ.Poller poller = CONTEXT.createPoller(1); - poller.register(subSocket, ZMQ.Poller.POLLIN); - Message testMessage = MessageBuilder.withPayload("test").setHeader("topic", "testTopic").build(); - messageHandler.handleMessage(testMessage).subscribe(); - while (true) { - poller.poll(10000); - if (poller.pollin(0)) { - ZMsg msg = ZMsg.recvMsg(subSocket); - assertThat(msg).isNotNull(); - assertThat(msg.unwrap().getString(ZMQ.CHARSET)).isEqualTo("testTopic"); - Message capturedMessage = new EmbeddedJsonHeadersMessageMapper().toMessage(msg.getFirst().getData()); - assertThat(capturedMessage).isEqualTo(testMessage); - msg.destroy(); - break; - } - } + await().untilAsserted(() -> { + messageHandler.handleMessage(testMessage).subscribe(); + ZMsg msg = ZMsg.recvMsg(subSocket); + assertThat(msg).isNotNull(); + assertThat(msg.unwrap().getString(ZMQ.CHARSET)).isEqualTo("testTopic"); + Message capturedMessage = new EmbeddedJsonHeadersMessageMapper().toMessage(msg.getFirst().getData()); + assertThat(capturedMessage).isEqualTo(testMessage); + msg.destroy(); + }); - poller.unregister(subSocket); - poller.close(); messageHandler.destroy(); subSocket.close(); }