Fix ZMQMessageHandlerTests race condition

Use `receiveTimeOut = 0` on SUB socket and rely
on the retries from the `await().atMost()`.
Also resubscribe in the retry callback before publishing the next attempt
This commit is contained in:
Artem Bilan
2021-07-09 11:38:03 -04:00
parent 0297f6bc9f
commit 39c3acb042

View File

@@ -86,10 +86,11 @@ public class ZeroMqMessageHandlerTests {
}
@Test
void testMessageHandlerForPubSub() throws InterruptedException {
void testMessageHandlerForPubSub() {
ZMQ.Socket subSocket = CONTEXT.createSocket(SocketType.SUB);
subSocket.setReceiveTimeOut(20_000);
subSocket.setReceiveTimeOut(0);
int port = subSocket.bindToRandomPort("tcp://*");
subSocket.subscribe("test");
ZeroMqMessageHandler messageHandler =
new ZeroMqMessageHandler(CONTEXT, "tcp://localhost:" + port, SocketType.PUB);
@@ -98,15 +99,12 @@ public class ZeroMqMessageHandlerTests {
new FunctionExpression<Message<?>>((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();