ZeroMQ test: Add sleep between SUB & PUB
This commit is contained in:
@@ -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) -> 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();
|
||||
|
||||
Reference in New Issue
Block a user