diff --git a/spring-integration-core/src/main/java/org/springframework/integration/events/IntegrationEvent.java b/spring-integration-core/src/main/java/org/springframework/integration/events/IntegrationEvent.java index 46ba1b8f9c..a56316eacb 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/events/IntegrationEvent.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/events/IntegrationEvent.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2019 the original author or authors. + * Copyright 2013-2020 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. @@ -48,6 +48,18 @@ public abstract class IntegrationEvent extends ApplicationEvent { return this.cause; } + /** + * Get the source as a specific type; the receiving variable must be declared with the + * correct type. + * @param the type. + * @return the source. + * @since 5.4 + */ + @SuppressWarnings("unchecked") + public T getSourceAsType() { + return (T) getSource(); + } + @Override public String toString() { return this.getClass().getSimpleName() + " [source=" + this.getSource() + diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/MqttComponent.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/MqttComponent.java new file mode 100644 index 0000000000..06bca8f7e3 --- /dev/null +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/MqttComponent.java @@ -0,0 +1,44 @@ +/* + * Copyright 2020 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.mqtt.core; + +import org.springframework.beans.factory.BeanNameAware; + +/** + * A component that interfaces with MQTT. + * + * @param The connection information type. + * + * @author Gary Russell + * @since 2.5 + * + */ +public interface MqttComponent extends BeanNameAware { + + /** + * Return this component's bean name. + * @return the bean name. + */ + String getBeanName(); + + /** + * Return information about the connection. + * @return the information. + */ + T getConnectionInfo(); + +} diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/MqttPahoComponent.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/MqttPahoComponent.java new file mode 100644 index 0000000000..482298c1fd --- /dev/null +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/MqttPahoComponent.java @@ -0,0 +1,33 @@ +/* + * Copyright 2020 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.mqtt.core; + +import org.eclipse.paho.client.mqttv3.MqttConnectOptions; + +/** + * An extension of {@link MqttComponent} for Eclipse Paho components. + * + * @author Gary Russell + * @since 5.4 + * + */ +public interface MqttPahoComponent extends MqttComponent { + + @Override + MqttConnectOptions getConnectionInfo(); + +} diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/event/MqttIntegrationEvent.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/event/MqttIntegrationEvent.java index 73dd83282b..b8cb024774 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/event/MqttIntegrationEvent.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/event/MqttIntegrationEvent.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2019 the original author or authors. + * Copyright 2014-2020 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. @@ -19,7 +19,9 @@ package org.springframework.integration.mqtt.event; import org.springframework.integration.events.IntegrationEvent; /** - * Base class for Mqtt Events. + * Base class for Mqtt Events. For {@link #getSourceAsType()}, you should use a sub type + * of {@link org.springframework.integration.mqtt.core.MqttComponent} for the receiving + * variable. * @author Gary Russell * * @since 4.1 diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/AbstractMqttMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/AbstractMqttMessageDrivenChannelAdapter.java index 23588c40c2..366b36a21d 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/AbstractMqttMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/AbstractMqttMessageDrivenChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2019 the original author or authors. + * Copyright 2002-2020 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. @@ -28,6 +28,7 @@ import org.springframework.integration.support.management.IntegrationManagedReso import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.jmx.export.annotation.ManagedOperation; import org.springframework.jmx.export.annotation.ManagedResource; +import org.springframework.lang.Nullable; import org.springframework.messaging.MessagingException; import org.springframework.util.Assert; @@ -54,7 +55,7 @@ public abstract class AbstractMqttMessageDrivenChannelAdapter extends MessagePro protected final Lock topicLock = new ReentrantLock(); // NOSONAR - public AbstractMqttMessageDrivenChannelAdapter(String url, String clientId, String... topic) { + public AbstractMqttMessageDrivenChannelAdapter(@Nullable String url, String clientId, String... topic) { Assert.hasText(clientId, "'clientId' cannot be null or empty"); Assert.notNull(topic, "'topics' cannot be null"); Assert.noNullElements(topic, "'topics' cannot have null elements"); @@ -110,6 +111,7 @@ public abstract class AbstractMqttMessageDrivenChannelAdapter extends MessagePro } } + @Nullable protected String getUrl() { return this.url; } diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java index 23e69a6906..570c921c05 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java @@ -35,15 +35,22 @@ import org.springframework.integration.acks.SimpleAcknowledgment; import org.springframework.integration.mqtt.core.ConsumerStopAction; import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory; import org.springframework.integration.mqtt.core.MqttPahoClientFactory; +import org.springframework.integration.mqtt.core.MqttPahoComponent; import org.springframework.integration.mqtt.event.MqttConnectionFailedEvent; +import org.springframework.integration.mqtt.event.MqttIntegrationEvent; import org.springframework.integration.mqtt.event.MqttSubscribedEvent; +import org.springframework.integration.mqtt.support.MqttUtils; import org.springframework.integration.support.AbstractIntegrationMessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessagingException; import org.springframework.util.Assert; /** - * Eclipse Paho Implementation. + * Eclipse Paho Implementation. When consuming {@link MqttIntegrationEvent}s published by + * this component use {@code MqttPahoComponent adapter = event.getSourceAsType()} to get a + * reference, allowing you to obtain the bean name and {@link MqttConnectOptions}. This + * technique allows consumption of events from both inbound and outbound endpoints in the + * same event listener. * * @author Gary Russell * @author Artem Bilan @@ -52,7 +59,7 @@ import org.springframework.util.Assert; * */ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDrivenChannelAdapter - implements MqttCallback, ApplicationEventPublisherAware { + implements MqttCallback, MqttPahoComponent, ApplicationEventPublisherAware { /** * The default completion timeout in milliseconds. @@ -89,9 +96,9 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv private ApplicationEventPublisher applicationEventPublisher; /** - * Use this constructor for a single url (although it may be overridden - * if the server URI(s) are provided by the {@link MqttConnectOptions#getServerURIs()} - * provided by the {@link MqttPahoClientFactory}). + * Use this constructor for a single url (although it may be overridden if the server + * URI(s) are provided by the {@link MqttConnectOptions#getServerURIs()} provided by + * the {@link MqttPahoClientFactory}). * @param url the URL. * @param clientId The client id. * @param clientFactory The client factory. @@ -99,13 +106,15 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv */ public MqttPahoMessageDrivenChannelAdapter(String url, String clientId, MqttPahoClientFactory clientFactory, String... topic) { + super(url, clientId, topic); this.clientFactory = clientFactory; } /** - * Use this constructor if the server URI(s) are provided by the {@link MqttConnectOptions#getServerURIs()} - * provided by the {@link MqttPahoClientFactory}. + * Use this constructor if the server URI(s) are provided by the + * {@link MqttConnectOptions#getServerURIs()} provided by the + * {@link MqttPahoClientFactory}. * @param clientId The client id. * @param clientFactory The client factory. * @param topic The topic(s). @@ -117,8 +126,9 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv this.clientFactory = clientFactory; } + /** - * Use this URL when you don't need additional {@link MqttConnectOptions}. + * Use this constructor when you don't need additional {@link MqttConnectOptions}. * @param url The URL. * @param clientId The client id. * @param topic The topic(s). @@ -174,6 +184,19 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv this.applicationEventPublisher = applicationEventPublisher; // NOSONAR (inconsistent synchronization) } + @Override + public MqttConnectOptions getConnectionInfo() { + MqttConnectOptions options = this.clientFactory.getConnectionOptions(); + if (options.getServerURIs() == null) { + String url = getUrl(); + if (url != null) { + options = MqttUtils.cloneConnectOptions(options); + options.setServerURIs(new String[] { url }); + } + } + return options; + } + @Override protected void doStart() { Assert.state(getTaskScheduler() != null, "A 'taskScheduler' is required"); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java index d50a79f481..4bdc6ee70e 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java @@ -27,6 +27,7 @@ import org.springframework.integration.handler.MessageProcessor; import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter; import org.springframework.integration.mqtt.support.MqttHeaders; import org.springframework.integration.mqtt.support.MqttMessageConverter; +import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.converter.MessageConverter; import org.springframework.util.Assert; @@ -67,7 +68,7 @@ public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler private int clientInstance; - public AbstractMqttMessageHandler(String url, String clientId) { + public AbstractMqttMessageHandler(@Nullable String url, String clientId) { Assert.hasText(clientId, "'clientId' cannot be null or empty"); this.url = url; this.clientId = clientId; @@ -185,6 +186,7 @@ public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler return this.converter; } + @Nullable protected String getUrl() { return this.url; } diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java index 73831583be..e291708444 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java @@ -27,17 +27,24 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory; import org.springframework.integration.mqtt.core.MqttPahoClientFactory; +import org.springframework.integration.mqtt.core.MqttPahoComponent; import org.springframework.integration.mqtt.event.MqttConnectionFailedEvent; +import org.springframework.integration.mqtt.event.MqttIntegrationEvent; import org.springframework.integration.mqtt.event.MqttMessageDeliveredEvent; import org.springframework.integration.mqtt.event.MqttMessageSentEvent; import org.springframework.integration.mqtt.support.MqttMessageConverter; +import org.springframework.integration.mqtt.support.MqttUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessagingException; import org.springframework.util.Assert; /** - * Eclipse Paho implementation. + * Eclipse Paho Implementation. When consuming {@link MqttIntegrationEvent}s published by + * this component use {@code MqttPahoComponent handler = event.getSourceAsType()} to get a + * reference, allowing you to obtain the bean name and {@link MqttConnectOptions}. This + * technique allows consumption of events from both inbound and outbound endpoints in the + * same event listener. * * @author Gary Russell * @author Artem Bilan @@ -46,7 +53,7 @@ import org.springframework.util.Assert; * */ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler - implements MqttCallback, ApplicationEventPublisherAware { + implements MqttCallback, MqttPahoComponent, ApplicationEventPublisherAware { /** * The default completion timeout in milliseconds. @@ -73,9 +80,9 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler private volatile IMqttAsyncClient client; /** - * Use this constructor for a single url (although it may be overridden - * if the server URI(s) are provided by the {@link MqttConnectOptions#getServerURIs()} - * provided by the {@link MqttPahoClientFactory}). + * Use this constructor for a single url (although it may be overridden if the server + * URI(s) are provided by the {@link MqttConnectOptions#getServerURIs()} provided by + * the {@link MqttPahoClientFactory}). * @param url the URL. * @param clientId The client id. * @param clientFactory The client factory. @@ -98,7 +105,7 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler } /** - * Use this URL when you don't need additional {@link MqttConnectOptions}. + * Use this constructor when you don't need additional {@link MqttConnectOptions}. * @param url The URL. * @param clientId The client id. */ @@ -154,6 +161,19 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler this.applicationEventPublisher = applicationEventPublisher; } + @Override + public MqttConnectOptions getConnectionInfo() { + MqttConnectOptions options = this.clientFactory.getConnectionOptions(); + if (options.getServerURIs() == null) { + String url = getUrl(); + if (url != null) { + options = MqttUtils.cloneConnectOptions(options); + options.setServerURIs(new String[] { url }); + } + } + return options; + } + @Override protected void onInit() { super.onInit(); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttUtils.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttUtils.java new file mode 100644 index 0000000000..22e67499bc --- /dev/null +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttUtils.java @@ -0,0 +1,49 @@ +/* + * Copyright 2020 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.mqtt.support; + +import org.eclipse.paho.client.mqttv3.MqttConnectOptions; + +import org.springframework.beans.BeanUtils; + +/** + * MQTT Utilities + * + * @author Gary Russell + * @since 5.4 + * + */ +public final class MqttUtils { + + private MqttUtils() { + } + + /** + * Clone the {@link MqttConnectOptions}, except the serverUris. + * @param options the options to clone. + * @return the clone. + */ + public static MqttConnectOptions cloneConnectOptions(MqttConnectOptions options) { + MqttConnectOptions options2 = new MqttConnectOptions(); + BeanUtils.copyProperties(options, options2, "password", "serverURIs"); + if (options.getPassword() != null) { + options2.setPassword(options.getPassword()); + } + return options2; + } + +} diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests-context.xml b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests-context.xml index cd6da5d918..258229c5e1 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests-context.xml +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests-context.xml @@ -27,4 +27,6 @@ + + diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests.java index 62c1a695f1..f7772b0008 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackToBackAdapterTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2019 the original author or authors. + * Copyright 2002-2020 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. @@ -20,6 +20,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; import static org.mockito.Mockito.mock; +import java.util.Arrays; import java.util.Collections; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -35,10 +36,13 @@ import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationEvent; import org.springframework.context.ApplicationEventPublisher; +import org.springframework.context.ApplicationListener; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory; +import org.springframework.integration.mqtt.core.MqttPahoComponent; import org.springframework.integration.mqtt.event.MqttMessageDeliveredEvent; import org.springframework.integration.mqtt.event.MqttMessageSentEvent; +import org.springframework.integration.mqtt.event.MqttSubscribedEvent; import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter; import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler; import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter; @@ -78,6 +82,9 @@ public class BackToBackAdapterTests { @Autowired private PollableChannel in; + @Autowired + private EventsListener listener; + @Test public void testSingleTopic() { MqttPahoMessageHandler adapter = new MqttPahoMessageHandler("tcp://localhost:1883", "si-test-out"); @@ -102,6 +109,7 @@ public class BackToBackAdapterTests { inbound.stop(); assertThat(out.getPayload()).isEqualTo("foo"); assertThat(out.getHeaders().get(MqttHeaders.RECEIVED_TOPIC)).isEqualTo("mqtt-foo"); + assertThat(adapter.getConnectionInfo().getServerURIs()[0]).isEqualTo("tcp://localhost:1883"); } @Test @@ -369,6 +377,20 @@ public class BackToBackAdapterTests { Message message = in.receive(20000); assertThat(message).isNotNull(); assertThat(message.getPayload()).isEqualTo("foo"); + MqttPahoComponent source = this.listener.event.getSourceAsType(); + assertThat(Arrays.toString(source.getConnectionInfo().getServerURIs())) + .isEqualTo("[tcp://localhost:1883, tcp://localhost:1883]"); + } + + public static class EventsListener implements ApplicationListener { + + volatile MqttSubscribedEvent event; + + @Override + public void onApplicationEvent(MqttSubscribedEvent event) { + this.event = event; + } + } private class EventPublisher implements ApplicationEventPublisher { diff --git a/src/reference/asciidoc/mqtt.adoc b/src/reference/asciidoc/mqtt.adoc index cb46dfa1a8..e52042bc32 100644 --- a/src/reference/asciidoc/mqtt.adoc +++ b/src/reference/asciidoc/mqtt.adoc @@ -403,3 +403,14 @@ Certain application events are published by the adapters. * `MqttSubscribedEvent` - published by the inbound adapter after subscribing to the topics. These events can be received by an `ApplicationListener` or with an `@EventListener` method. + +To determine the source of an event, use the following; you can check the bean name and/or the connect options (to access the server URIs etc). + +==== +[source, java] +---- +MqttPahoComponent source = event.getSourceAsType(); +String beanName = source.getBeanName(); +MqttConnectOptions options = source.getConnectionInfo(); +---- +====