GH-183: Expose MQTT SSL properties
Fixes https://github.com/spring-cloud/stream-applications/issues/183 * Add `MqttProperties.sslProperties` common MQTT configuration property as a `Map` structure * Delegate a content fo the `MqttProperties.sslProperties` into `mqttConnectOptions.setSSLProperties()` in the common `MqttConfiguration` * Verify that SSL properties are applied in both `MqttSupplierTests` & `MqttConsumerTests` * Document a new `MqttProperties.sslProperties` configuration property for both consumer and supplier * Use `eclipse-mosquitto` Docker image for MQTT tests instead of RabbitMQ plugin: the Mosquitto starts much faster
This commit is contained in:
@@ -11,7 +11,7 @@
|
||||
|
||||
<artifactId>mqtt-common</artifactId>
|
||||
<name>mqtt-common</name>
|
||||
<description>file consumer</description>
|
||||
<description>MQTT Common</description>
|
||||
|
||||
<dependencies>
|
||||
|
||||
|
||||
@@ -16,6 +16,9 @@
|
||||
|
||||
package org.springframework.cloud.fn.common.mqtt;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.Properties;
|
||||
|
||||
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
|
||||
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
|
||||
import org.eclipse.paho.client.mqttv3.persist.MqttDefaultFilePersistence;
|
||||
@@ -31,7 +34,7 @@ import org.springframework.util.ObjectUtils;
|
||||
* Generic mqtt configuration.
|
||||
*
|
||||
* @author Janne Valkealahti
|
||||
*
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@Configuration
|
||||
public class MqttConfiguration {
|
||||
@@ -50,6 +53,14 @@ public class MqttConfiguration {
|
||||
mqttConnectOptions.setConnectionTimeout(mqttProperties.getConnectionTimeout());
|
||||
mqttConnectOptions.setKeepAliveInterval(mqttProperties.getKeepAliveInterval());
|
||||
|
||||
Map<String, String> sslProperties = mqttProperties.getSslProperties();
|
||||
|
||||
if (!sslProperties.isEmpty()) {
|
||||
Properties sslProps = new Properties();
|
||||
sslProps.putAll(sslProperties);
|
||||
mqttConnectOptions.setSSLProperties(sslProps);
|
||||
}
|
||||
|
||||
DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
|
||||
factory.setConnectionOptions(mqttConnectOptions);
|
||||
|
||||
@@ -61,4 +72,5 @@ public class MqttConfiguration {
|
||||
}
|
||||
return factory;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,6 +16,9 @@
|
||||
|
||||
package org.springframework.cloud.fn.common.mqtt;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import javax.validation.constraints.Size;
|
||||
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
@@ -25,7 +28,7 @@ import org.springframework.validation.annotation.Validated;
|
||||
* Generic mqtt connection properties.
|
||||
*
|
||||
* @author Janne Valkealahti
|
||||
*
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@Validated
|
||||
@ConfigurationProperties("mqtt")
|
||||
@@ -71,6 +74,11 @@ public class MqttProperties {
|
||||
*/
|
||||
private String persistenceDirectory = "/tmp/paho";
|
||||
|
||||
/**
|
||||
* MQTT Client SSL properties.
|
||||
*/
|
||||
private final Map<String, String> sslProperties = new HashMap<>();
|
||||
|
||||
@Size(min = 1)
|
||||
public String[] getUrl() {
|
||||
return url;
|
||||
@@ -135,4 +143,9 @@ public class MqttProperties {
|
||||
public void setPersistenceDirectory(String persistenceDirectory) {
|
||||
this.persistenceDirectory = persistenceDirectory;
|
||||
}
|
||||
|
||||
public Map<String, String> getSslProperties() {
|
||||
return this.sslProperties;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,6 +16,12 @@ All configuration properties are prefixed with `mqtt.consumer`.
|
||||
|
||||
For more information on the various options available, please see link:src/main/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerProperties.java[MqttConsumerProperties].
|
||||
|
||||
## SSL Configuration
|
||||
|
||||
The MQTT Paho client can accept an SSL configuration via `MqttConnectOptions.setSSLProperties()`.
|
||||
These properties are exposed on the `MqttProperties.sslProperties` map.
|
||||
The keys for these SSL properties should be taken from the `org.eclipse.paho.client.mqttv3.internal.security.SSLSocketFactoryFactory` constants, which all start with the `com.ibm.ssl.` prefix.
|
||||
|
||||
## Tests
|
||||
|
||||
See this link:src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java[test suite] for the various ways, this consumer is used.
|
||||
|
||||
@@ -16,8 +16,11 @@
|
||||
|
||||
package org.springframework.cloud.fn.consumer.mqtt;
|
||||
|
||||
import java.util.Properties;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
|
||||
import org.eclipse.paho.client.mqttv3.internal.security.SSLSocketFactoryFactory;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.Tag;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -32,25 +35,35 @@ import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
|
||||
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
|
||||
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
@SpringBootTest(properties = "mqtt.consumer.topic=test")
|
||||
@SpringBootTest(properties = {
|
||||
"mqtt.consumer.topic=test",
|
||||
"mqtt.ssl-properties.com.ibm.ssl.protocol=TLS",
|
||||
"mqtt.ssl-properties.com.ibm.ssl.keyStoreType=TEST" })
|
||||
@DirtiesContext
|
||||
@Tag("integration")
|
||||
public class MqttConsumerTests {
|
||||
|
||||
static {
|
||||
GenericContainer mosquitto = new GenericContainer("cyrilix/rabbitmq-mqtt")
|
||||
.withExposedPorts(1883);
|
||||
GenericContainer<?> mosquitto =
|
||||
new GenericContainer<>("eclipse-mosquitto:2.0.13")
|
||||
.withCommand("mosquitto -c /mosquitto-no-auth.conf")
|
||||
.withReuse(true)
|
||||
.withExposedPorts(1883);
|
||||
mosquitto.start();
|
||||
final Integer mappedPort = mosquitto.getMappedPort(1883);
|
||||
System.setProperty("mqtt.url", "tcp://localhost:" + mappedPort);
|
||||
}
|
||||
|
||||
@Autowired
|
||||
private MqttPahoMessageDrivenChannelAdapter mqttPahoMessageDrivenChannelAdapter;
|
||||
|
||||
@Autowired
|
||||
private Consumer<Message<?>> mqttConsumer;
|
||||
|
||||
@@ -64,6 +77,12 @@ public class MqttConsumerTests {
|
||||
|
||||
@Test
|
||||
public void testMqttConsumer() {
|
||||
MqttConnectOptions connectionInfo = this.mqttPahoMessageDrivenChannelAdapter.getConnectionInfo();
|
||||
Properties sslProperties = connectionInfo.getSSLProperties();
|
||||
assertThat(sslProperties)
|
||||
.containsEntry(SSLSocketFactoryFactory.SSLPROTOCOL, SSLSocketFactoryFactory.DEFAULT_PROTOCOL)
|
||||
.containsEntry(SSLSocketFactoryFactory.KEYSTORETYPE, "TEST");
|
||||
|
||||
this.mqttConsumer.accept(MessageBuilder.withPayload("hello").build());
|
||||
Message<?> in = this.queue.receive(10000);
|
||||
assertThat(in).isNotNull();
|
||||
|
||||
@@ -24,6 +24,12 @@ All configuration properties are prefixed with `mqtt.supplier` and `mqtt`.
|
||||
|
||||
For more information on the various options available, please see link:src/main/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierProperties.java[MqttSupplierProperties].
|
||||
|
||||
## SSL Configuration
|
||||
|
||||
The MQTT Paho client can accept an SSL configuration via `MqttConnectOptions.setSSLProperties()`.
|
||||
These properties are exposed on the `MqttProperties.sslProperties` map.
|
||||
The keys for these SSL properties should be taken from the `org.eclipse.paho.client.mqttv3.internal.security.SSLSocketFactoryFactory` constants, which all start with the `com.ibm.ssl.` prefix.
|
||||
|
||||
## Tests
|
||||
|
||||
See this link:src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java[test suite] for the various ways, this supplier is used.
|
||||
|
||||
@@ -16,14 +16,17 @@
|
||||
|
||||
package org.springframework.cloud.fn.supplier.mqtt;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
import java.util.Properties;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
|
||||
import org.eclipse.paho.client.mqttv3.internal.security.SSLSocketFactoryFactory;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.Tag;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.testcontainers.containers.GenericContainer;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
@@ -37,7 +40,8 @@ import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
/**
|
||||
* Tests for Mqtt Supplier.
|
||||
@@ -48,14 +52,21 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
* @author Artem Bilan
|
||||
*
|
||||
*/
|
||||
@SpringBootTest(properties = {"mqtt.supplier.topics=test,fake", "mqtt.supplier.qos=0,0"})
|
||||
@SpringBootTest(properties =
|
||||
{ "mqtt.supplier.topics=test,fake",
|
||||
"mqtt.supplier.qos=0,0",
|
||||
"mqtt.ssl-properties.com.ibm.ssl.protocol=TLS",
|
||||
"mqtt.ssl-properties.com.ibm.ssl.keyStoreType=TEST"})
|
||||
@DirtiesContext
|
||||
@Tag("integration")
|
||||
public class MqttSupplierTests {
|
||||
|
||||
static {
|
||||
GenericContainer mosquitto = new GenericContainer("cyrilix/rabbitmq-mqtt")
|
||||
.withExposedPorts(1883);
|
||||
GenericContainer<?> mosquitto =
|
||||
new GenericContainer<>("eclipse-mosquitto:2.0.13")
|
||||
.withCommand("mosquitto -c /mosquitto-no-auth.conf")
|
||||
.withReuse(true)
|
||||
.withExposedPorts(1883);
|
||||
mosquitto.start();
|
||||
final Integer mappedPort = mosquitto.getMappedPort(1883);
|
||||
System.setProperty("mqtt.url", "tcp://localhost:" + mappedPort);
|
||||
@@ -70,12 +81,16 @@ public class MqttSupplierTests {
|
||||
private Supplier<Flux<Message<?>>> mqttSupplier;
|
||||
|
||||
@Autowired
|
||||
private MessageHandler mqttOutbound;
|
||||
private MqttPahoMessageHandler mqttOutbound;
|
||||
|
||||
@Test
|
||||
public void testBasicFlow() {
|
||||
|
||||
mqttOutbound.handleMessage(MessageBuilder.withPayload("hello").build());
|
||||
MqttConnectOptions connectionInfo = this.mqttOutbound.getConnectionInfo();
|
||||
Properties sslProperties = connectionInfo.getSSLProperties();
|
||||
assertThat(sslProperties)
|
||||
.containsEntry(SSLSocketFactoryFactory.SSLPROTOCOL, SSLSocketFactoryFactory.DEFAULT_PROTOCOL)
|
||||
.containsEntry(SSLSocketFactoryFactory.KEYSTORETYPE, "TEST");
|
||||
this.mqttOutbound.handleMessage(MessageBuilder.withPayload("hello").build());
|
||||
|
||||
final Flux<Message<?>> messageFlux = mqttSupplier.get();
|
||||
|
||||
@@ -108,5 +123,7 @@ public class MqttSupplierTests {
|
||||
public DefaultPahoMessageConverter producerConverter() {
|
||||
return new DefaultPahoMessageConverter(1, true, "UTF-8");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user