GH-9228: Provide binding to ZeroMqMessageHandler

Fixes: #9228

* add docs

* protected constructor in ZeroMqMessageHandlerSpec, expose them via ZeroMq

* introduce ZeroMqUtils, for common Zero MQ utilities functions

* use ZeroMqUtils.bindSocket in ZeroMqMessageProducer

* refactor ZeroMqMessageHandler providing connectUrl and bindPort setters and simple constructors, following the same logic used for ZeroMqMessageProvider

* fix tests to follow the new ZeroMqMessageHandler implementation

* fix typo

* add new updates to whats-new.adoc

* address ZeroMQUtils comments

* remove connectUrl and boundPort setters, add Javadoc to new constructors

* add new DLS constructor for random port

* add since closure in ZeroMqUtils

* add DSL support methods and specific that, when not defined, the socket will be bound to a random port
This commit is contained in:
Alessio Matricardi
2024-07-01 19:00:26 +02:00
committed by GitHub
parent 78ae221fd5
commit 272dde8945
9 changed files with 302 additions and 33 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2020 the original author or authors.
* Copyright 2020-2024 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -17,7 +17,7 @@
package org.springframework.integration.zeromq;
/**
* The message headers constants to repsent ZeroMq message attributes.
* The message headers constants to represent ZeroMq message attributes.
*
* @author Artem Bilan
*

View File

@@ -0,0 +1,53 @@
/*
* Copyright 2020-2024 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.zeromq;
import org.zeromq.ZMQ;
/**
* Module that wraps common methods of ZeroMq integration classes
*
* @author Alessio Matricardi
*
* @since 6.4
*
*/
public final class ZeroMqUtils {
/**
* Bind the ZeroMq socket to the given port over the TCP transport protocol.
* @param socket the ZeroMq socket
* @param port the port to bind ZeroMq socket to over TCP. If equal to 0, the socket will bind to a random port.
* @return the effectively bound port
*/
public static int bindSocket(ZMQ.Socket socket, int port) {
if (port == 0) {
return socket.bindToRandomPort("tcp://*");
}
else {
boolean bound = socket.bind("tcp://*:" + port);
if (!bound) {
throw new IllegalArgumentException("Cannot bind ZeroMQ socket to port: " + port);
}
return port;
}
}
private ZeroMqUtils() {
}
}

View File

@@ -25,6 +25,7 @@ import org.zeromq.ZContext;
* Factory class for ZeroMq components DSL.
*
* @author Artem Bilan
* @author Alessio Matricardi
*
* @since 5.4
*/
@@ -58,6 +59,17 @@ public final class ZeroMq {
return outboundChannelAdapter(context, () -> connectUrl);
}
/**
* Create an instance of {@link ZeroMqMessageHandlerSpec} for the provided {@link ZContext} and binding port.
* @param context the {@link ZContext} to use.
* @param port the port to bind ZeroMq socket to over TCP.
* @return the spec.
* @since 6.4
*/
public static ZeroMqMessageHandlerSpec outboundChannelAdapter(ZContext context, int port) {
return new ZeroMqMessageHandlerSpec(context, port);
}
/**
* Create an instance of {@link ZeroMqMessageHandlerSpec} for the provided {@link ZContext}
* and connection URL supplier.
@@ -84,6 +96,43 @@ public final class ZeroMq {
return new ZeroMqMessageHandlerSpec(context, connectUrl, socketType);
}
/**
* Create an instance of {@link ZeroMqMessageHandlerSpec} for the provided {@link ZContext}.
* The created socket will be bound to a random port.
* @param context the {@link ZContext} to use.
* @return the spec.
* @since 6.4
*/
public static ZeroMqMessageHandlerSpec outboundChannelAdapter(ZContext context) {
return new ZeroMqMessageHandlerSpec(context);
}
/**
* Create an instance of {@link ZeroMqMessageHandlerSpec} for the provided {@link ZContext} and {@link SocketType}.
* The created socket will be bound to a random port.
* @param context the {@link ZContext} to use.
* @param socketType the {@link SocketType} for ZeroMq socket.
* @return the spec.
* @since 6.4
*/
public static ZeroMqMessageHandlerSpec outboundChannelAdapter(ZContext context, SocketType socketType) {
return new ZeroMqMessageHandlerSpec(context, socketType);
}
/**
* Create an instance of {@link ZeroMqMessageHandlerSpec} for the provided {@link ZContext}, binding port
* and {@link SocketType}.
* @param context the {@link ZContext} to use.
* @param port the port to bind ZeroMq socket to over TCP.
* @param socketType the {@link SocketType} for ZeroMq socket.
* @return the spec.
* @since 6.4
*/
public static ZeroMqMessageHandlerSpec outboundChannelAdapter(ZContext context, int port,
SocketType socketType) {
return new ZeroMqMessageHandlerSpec(context, port, socketType);
}
/**
* Create an instance of {@link ZeroMqMessageHandlerSpec} for the provided {@link ZContext},
* connection URL supplier and {@link SocketType}.

View File

@@ -52,6 +52,26 @@ public class ZeroMqMessageHandlerSpec
this(context, () -> connectUrl);
}
/**
* Create an instance based on the provided {@link ZContext}.
* The created socket will be bound to a random port.
* @param context the {@link ZContext} to use for creating sockets.
* @since 6.4
*/
protected ZeroMqMessageHandlerSpec(ZContext context) {
this(context, SocketType.PAIR);
}
/**
* Create an instance based on the provided {@link ZContext} and binding port.
* @param context the {@link ZContext} to use for creating sockets.
* @param port the port to bind ZeroMq socket to over TCP.
* @since 6.4
*/
protected ZeroMqMessageHandlerSpec(ZContext context, int port) {
this(context, port, SocketType.PAIR);
}
/**
* Create an instance based on the provided {@link ZContext} and connection string supplier.
* @param context the {@link ZContext} to use for creating sockets.
@@ -73,6 +93,30 @@ public class ZeroMqMessageHandlerSpec
this(context, () -> connectUrl, socketType);
}
/**
* Create an instance based on the provided {@link ZContext} and {@link SocketType}.
* The created socket will be bound to a random port.
* @param context the {@link ZContext} to use for creating sockets.
* @param socketType the {@link SocketType} to use;
* only {@link SocketType#PAIR}, {@link SocketType#PUB} and {@link SocketType#PUSH} are supported.
* @since 6.4
*/
protected ZeroMqMessageHandlerSpec(ZContext context, SocketType socketType) {
super(new ZeroMqMessageHandler(context, socketType));
}
/**
* Create an instance based on the provided {@link ZContext}, binding port and {@link SocketType}.
* @param context the {@link ZContext} to use for creating sockets.
* @param port the port to bind ZeroMq socket to over TCP.
* @param socketType the {@link SocketType} to use;
* only {@link SocketType#PAIR}, {@link SocketType#PUB} and {@link SocketType#PUSH} are supported.
* @since 6.4
*/
protected ZeroMqMessageHandlerSpec(ZContext context, int port, SocketType socketType) {
super(new ZeroMqMessageHandler(context, port, socketType));
}
/**
* Create an instance based on the provided {@link ZContext}, connection string supplier and {@link SocketType}.
* @param context the {@link ZContext} to use for creating sockets.

View File

@@ -40,6 +40,7 @@ import org.springframework.integration.mapping.InboundMessageMapper;
import org.springframework.integration.support.converter.ConfigurableCompositeMessageConverter;
import org.springframework.integration.support.management.IntegrationManagedResource;
import org.springframework.integration.zeromq.ZeroMqHeaders;
import org.springframework.integration.zeromq.ZeroMqUtils;
import org.springframework.jmx.export.annotation.ManagedOperation;
import org.springframework.jmx.export.annotation.ManagedResource;
import org.springframework.lang.Nullable;
@@ -263,7 +264,7 @@ public class ZeroMqMessageProducer extends MessageProducerSupport {
socket.connect(this.connectUrl);
}
else {
this.bindPort.set(bindSocket(socket, this.bindPort.get()));
this.bindPort.set(ZeroMqUtils.bindSocket(socket, this.bindPort.get()));
}
})
.cache()
@@ -319,17 +320,4 @@ public class ZeroMqMessageProducer extends MessageProducerSupport {
this.socketMono.doOnNext(ZMQ.Socket::close).block();
}
private static int bindSocket(ZMQ.Socket socket, int port) {
if (port == 0) {
return socket.bindToRandomPort("tcp://*");
}
else {
boolean bound = socket.bind("tcp://*:" + port);
if (!bound) {
throw new IllegalArgumentException("Cannot bind ZeroMQ socket to port: " + port);
}
return port;
}
}
}

View File

@@ -18,6 +18,7 @@ package org.springframework.integration.zeromq.outbound;
import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
import java.util.function.Supplier;
@@ -43,6 +44,8 @@ import org.springframework.integration.mapping.ConvertingBytesMessageMapper;
import org.springframework.integration.mapping.OutboundMessageMapper;
import org.springframework.integration.support.converter.ConfigurableCompositeMessageConverter;
import org.springframework.integration.support.management.ManageableLifecycle;
import org.springframework.integration.zeromq.ZeroMqUtils;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.converter.MessageConverter;
import org.springframework.util.Assert;
@@ -50,14 +53,14 @@ import org.springframework.util.Assert;
/**
* The {@link AbstractReactiveMessageHandler} implementation for publishing messages over ZeroMq socket.
* Only {@link SocketType#PAIR}, {@link SocketType#PUB} and {@link SocketType#PUSH} are supported.
* This component is only connecting (no Binding) to another side, e.g. ZeroMq proxy.
* This component can bind or connect the socket.
* <p>
* When the {@link SocketType#PUB} is used, the {@link #topicExpression} is evaluated against a
* request message to inject a topic frame into a ZeroMq message if it is not {@code null}.
* The subscriber side must receive the topic frame first before parsing the actual data.
* <p>
* When the payload of the request message is a {@link ZMsg}, no any conversion and topic extraction happen:
* the {@link ZMsg} is sent into a socket as is and it is not destroyed for possible further reusing.
* the {@link ZMsg} is sent into a socket as is, and it is not destroyed for possible further reusing.
*
* @author Artem Bilan
* @author Alessio Matricardi
@@ -74,7 +77,7 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
private final Scheduler publisherScheduler = Schedulers.newSingle("zeroMqMessageHandlerScheduler");
private final Mono<ZMQ.Socket> socketMono;
private volatile Mono<ZMQ.Socket> socketMono;
private OutboundMessageMapper<byte[]> messageMapper;
@@ -91,6 +94,38 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
private volatile boolean wrapTopic = true;
private final ZContext context;
private final SocketType socketType;
private final AtomicInteger bindPort = new AtomicInteger();
@Nullable
private Supplier<String> connectUrl;
/**
* Create an instance based on the provided {@link ZContext}.
* @param context the {@link ZContext} to use for creating sockets.
* @since 6.4
*/
public ZeroMqMessageHandler(ZContext context) {
this(context, SocketType.PAIR);
}
/**
* Create an instance based on the provided {@link ZContext} and {@link SocketType}.
* @param context the {@link ZContext} to use for creating sockets.
* @param socketType the {@link SocketType} to use;
* only {@link SocketType#PAIR}, {@link SocketType#PUB} and {@link SocketType#PUSH} are supported.
*/
public ZeroMqMessageHandler(ZContext context, SocketType socketType) {
Assert.notNull(context, "'context' must not be null");
Assert.state(VALID_SOCKET_TYPES.contains(socketType),
() -> "'socketType' can only be one of the: " + VALID_SOCKET_TYPES);
this.context = context;
this.socketType = socketType;
}
/**
* Create an instance based on the provided {@link ZContext} and connection string.
* @param context the {@link ZContext} to use for creating sockets.
@@ -100,6 +135,16 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
this(context, connectUrl, SocketType.PAIR);
}
/**
* Create an instance based on the provided {@link ZContext} and binding port.
* @param context the {@link ZContext} to use for creating sockets.
* @param port the port to bind ZeroMq socket to over TCP.
* @since 6.4
*/
public ZeroMqMessageHandler(ZContext context, int port) {
this(context, port, SocketType.PAIR);
}
/**
* Create an instance based on the provided {@link ZContext} and connection string supplier.
* @param context the {@link ZContext} to use for creating sockets.
@@ -122,6 +167,20 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
Assert.hasText(connectUrl, "'connectUrl' must not be empty");
}
/**
* Create an instance based on the provided {@link ZContext}, binding port and {@link SocketType}.
* @param context the {@link ZContext} to use for creating sockets.
* @param port the port to bind ZeroMq socket to over TCP.
* @param socketType the {@link SocketType} to use;
* only {@link SocketType#PAIR}, {@link SocketType#PUB} and {@link SocketType#PUSH} are supported.
* @since 6.4
*/
public ZeroMqMessageHandler(ZContext context, int port, SocketType socketType) {
this(context, socketType);
Assert.isTrue(port > 0, "'port' must not be zero or negative");
this.bindPort.set(port);
}
/**
* Create an instance based on the provided {@link ZContext}, connection string supplier and {@link SocketType}.
* @param context the {@link ZContext} to use for creating sockets.
@@ -131,17 +190,9 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
* @since 5.5.9
*/
public ZeroMqMessageHandler(ZContext context, Supplier<String> connectUrl, SocketType socketType) {
Assert.notNull(context, "'context' must not be null");
this(context, socketType);
Assert.notNull(connectUrl, "'connectUrl' must not be null");
Assert.state(VALID_SOCKET_TYPES.contains(socketType),
() -> "'socketType' can only be one of the: " + VALID_SOCKET_TYPES);
this.socketMono =
Mono.just(context.createSocket(socketType))
.publishOn(this.publisherScheduler)
.doOnNext((socket) -> this.socketConfigurer.accept(socket))
.doOnNext((socket) -> socket.connect(connectUrl.get()))
.cache()
.publishOn(this.publisherScheduler);
this.connectUrl = connectUrl;
}
/**
@@ -206,6 +257,16 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
this.wrapTopic = wrapTopic;
}
/**
* Return the port a socket is bound or 0 if this message producer has not been started yet
* or the socket is connected - not bound.
* @return the port for a socket or 0.
* @since 6.4
*/
public int getBoundPort() {
return this.bindPort.get();
}
@Override
public String getComponentType() {
return "zeromq:outbound-channel-adapter";
@@ -228,6 +289,20 @@ public class ZeroMqMessageHandler extends AbstractReactiveMessageHandler
@Override
public void start() {
if (!this.running.getAndSet(true)) {
this.socketMono =
Mono.just(this.context.createSocket(this.socketType))
.publishOn(this.publisherScheduler)
.doOnNext((socket) -> this.socketConfigurer.accept(socket))
.doOnNext((socket) -> {
if (this.connectUrl != null) {
socket.connect(this.connectUrl.get());
}
else {
this.bindPort.set(ZeroMqUtils.bindSocket(socket, this.bindPort.get()));
}
})
.cache()
.publishOn(this.publisherScheduler);
this.socketMonoSubscriber = this.socketMono.subscribe();
}
}

View File

@@ -35,6 +35,7 @@ import org.springframework.integration.zeromq.ZeroMqProxy;
import org.springframework.messaging.Message;
import org.springframework.messaging.converter.ByteArrayMessageConverter;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.util.TestSocketUtils;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
@@ -65,12 +66,12 @@ public class ZeroMqMessageHandlerTests {
messageHandler.setBeanFactory(mock(BeanFactory.class));
messageHandler.setSocketConfigurer(s -> s.setZapDomain("global"));
messageHandler.afterPropertiesSet();
messageHandler.start();
@SuppressWarnings("unchecked")
Mono<ZMQ.Socket> socketMono = TestUtils.getPropertyValue(messageHandler, "socketMono", Mono.class);
ZMQ.Socket socketInUse = socketMono.block(Duration.ofSeconds(10));
assertThat(socketInUse.getZapDomain()).isEqualTo("global");
messageHandler.start();
Message<?> testMessage = new GenericMessage<>("test");
messageHandler.handleMessage(testMessage).subscribe();
@@ -187,4 +188,41 @@ public class ZeroMqMessageHandlerTests {
subSocket.close();
}
@Test
void testMessageHandlerForPubSubWithBind() {
int boundPort = TestSocketUtils.findAvailableTcpPort();
ZeroMqMessageHandler messageHandler =
new ZeroMqMessageHandler(CONTEXT, boundPort, SocketType.PUB);
messageHandler.setBeanFactory(mock(BeanFactory.class));
messageHandler.setTopicExpression(
new FunctionExpression<Message<?>>((message) -> message.getHeaders().get("topic")));
messageHandler.setMessageMapper(new EmbeddedJsonHeadersMessageMapper());
messageHandler.wrapTopic(false);
messageHandler.afterPropertiesSet();
messageHandler.start();
ZMQ.Socket subSocket = CONTEXT.createSocket(SocketType.SUB);
subSocket.setReceiveTimeOut(0);
subSocket.connect("tcp://localhost:" + boundPort);
subSocket.subscribe("test");
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();
assertThat(msg.pop().getString(ZMQ.CHARSET)).isEqualTo("testTopic");
Message<?> capturedMessage =
new EmbeddedJsonHeadersMessageMapper().toMessage(msg.getFirst().getData());
assertThat(capturedMessage).isEqualTo(testMessage);
msg.destroy();
});
messageHandler.destroy();
subSocket.close();
}
}