diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java index dc59ac0dcb..8e70534ee5 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java @@ -28,14 +28,13 @@ import org.springframework.http.server.PathContainer; import org.springframework.http.server.RequestPath; import org.springframework.web.HttpRequestHandler; import org.springframework.web.servlet.HandlerExecutionChain; -import org.springframework.web.servlet.handler.AbstractHandlerMapping; import org.springframework.web.servlet.handler.AbstractUrlHandlerMapping; import org.springframework.web.util.ServletRequestPathUtils; import org.springframework.web.util.pattern.PathPattern; import org.springframework.web.util.pattern.PathPatternParser; /** - * The {@link AbstractHandlerMapping} implementation for dynamic WebSocket endpoint registrations in Spring Integration. + * The {@link AbstractUrlHandlerMapping} implementation for dynamic WebSocket endpoint registrations in Spring Integration. *
* TODO until https://github.com/spring-projects/spring-framework/issues/26798 * 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 8793c5b26c..296fdb496d 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() throws InterruptedException { + void testMessageHandlerForPubSub() { ZMQ.Socket subSocket = CONTEXT.createSocket(SocketType.SUB); subSocket.setReceiveTimeOut(20_000); int port = subSocket.bindToRandomPort("tcp://*"); @@ -100,19 +100,27 @@ public class ZeroMqMessageHandlerTests { messageHandler.setMessageMapper(new EmbeddedJsonHeadersMessageMapper()); messageHandler.afterPropertiesSet(); - // Give it some time to bind and subscribe - Thread.sleep(2000); + 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(); - 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); + 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; + } + } - msg.destroy(); + poller.unregister(subSocket); + poller.close(); messageHandler.destroy(); subSocket.close(); }