diff --git a/common/mqtt-common/pom.xml b/common/mqtt-common/pom.xml
index 8965daed..43c1af3b 100644
--- a/common/mqtt-common/pom.xml
+++ b/common/mqtt-common/pom.xml
@@ -11,7 +11,7 @@
mqtt-common
mqtt-common
- file consumer
+ MQTT Common
diff --git a/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttConfiguration.java b/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttConfiguration.java
index 0335940b..7cdad151 100644
--- a/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttConfiguration.java
+++ b/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttConfiguration.java
@@ -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 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;
}
+
}
diff --git a/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttProperties.java b/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttProperties.java
index e8c414f8..e704698d 100644
--- a/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttProperties.java
+++ b/common/mqtt-common/src/main/java/org/springframework/cloud/fn/common/mqtt/MqttProperties.java
@@ -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 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 getSslProperties() {
+ return this.sslProperties;
+ }
+
}
diff --git a/consumer/mqtt-consumer/README.adoc b/consumer/mqtt-consumer/README.adoc
index 0878f75b..5217496b 100644
--- a/consumer/mqtt-consumer/README.adoc
+++ b/consumer/mqtt-consumer/README.adoc
@@ -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.
diff --git a/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java b/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java
index 414f30d7..dad13103 100644
--- a/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java
+++ b/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java
@@ -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> 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();
diff --git a/supplier/mqtt-supplier/README.adoc b/supplier/mqtt-supplier/README.adoc
index 33b1a24c..3c6c3e99 100644
--- a/supplier/mqtt-supplier/README.adoc
+++ b/supplier/mqtt-supplier/README.adoc
@@ -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.
diff --git a/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java b/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java
index 4705f644..21c66ece 100644
--- a/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java
+++ b/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java
@@ -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>> 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> messageFlux = mqttSupplier.get();
@@ -108,5 +123,7 @@ public class MqttSupplierTests {
public DefaultPahoMessageConverter producerConverter() {
return new DefaultPahoMessageConverter(1, true, "UTF-8");
}
+
}
+
}