From d3924db5a16fbc071526c5455efc4e30f669bcea Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 30 Jan 2024 17:03:24 -0500 Subject: [PATCH] GH-8873: Fix on-demand subscription for MQTT v5 Fixes: #8873 The `mqttClient.subscribe()` API does not check if properties are provided and fails on the `subscriptionProperties.getSubscriptionIdentifiers().get(0)` call with an `IndexOutOfBoundsException` * Use another `mqttClient.subscribe()` API in the `Mqttv5PahoMessageDrivenChannelAdapter` where there is not such a check * Ensure that `addTopic(NAME)` works as expected in the `Mqttv5BackToBackTests` **Cherry-pick to `6.2.x` & `6.1.x`** --- .../Mqttv5PahoMessageDrivenChannelAdapter.java | 5 +++-- .../integration/mqtt/Mqttv5BackToBackTests.java | 17 +++++++++++++++++ 2 files changed, 20 insertions(+), 2 deletions(-) diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java index a51d0d4593..cdb59906e2 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java @@ -336,7 +336,8 @@ public class Mqttv5PahoMessageDrivenChannelAdapter this.subscriptions.add(subscription); } if (this.mqttClient != null && this.mqttClient.isConnected()) { - this.mqttClient.subscribe(subscription, this::messageArrived) + this.mqttClient.subscribe(new MqttSubscription[] { subscription }, + null, null, new IMqttMessageListener[] { this::messageArrived }, new MqttProperties()) .waitForCompletion(getCompletionTimeout()); } } @@ -466,7 +467,7 @@ public class Mqttv5PahoMessageDrivenChannelAdapter IMqttMessageListener[] listeners = IntStream.range(0, mqttSubscriptions.length) .mapToObj(t -> listener) .toArray(IMqttMessageListener[]::new); - this.mqttClient.subscribe(mqttSubscriptions, null, null, listeners, null) + this.mqttClient.subscribe(mqttSubscriptions, null, null, listeners, new MqttProperties()) .waitForCompletion(getCompletionTimeout()); String message = "Connected and subscribed to " + Arrays.toString(mqttSubscriptions); logger.debug(message); diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/Mqttv5BackToBackTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/Mqttv5BackToBackTests.java index aa1cba6e39..80e596d0cf 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/Mqttv5BackToBackTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/Mqttv5BackToBackTests.java @@ -75,6 +75,9 @@ public class Mqttv5BackToBackTests implements MosquittoContainerTest { @Autowired private Config config; + @Autowired + private Mqttv5PahoMessageDrivenChannelAdapter mqttv5MessageDrivenChannelAdapter; + @Test //GH-3732 public void testNoNpeIsNotThrownInCaseDoInitIsNotInvokedBeforeTopicAddition() { Mqttv5PahoMessageDrivenChannelAdapter channelAdapter = @@ -118,6 +121,20 @@ public class Mqttv5BackToBackTests implements MosquittoContainerTest { .hasAtLeastOneElementOfType(MqttMessageSentEvent.class) .hasAtLeastOneElementOfType(MqttMessageDeliveredEvent.class) .hasAtLeastOneElementOfType(MqttSubscribedEvent.class); + + this.mqttv5MessageDrivenChannelAdapter.addTopic("anotherTopic"); + + testPayload = "another payload"; + + this.mqttOutFlowInput.send( + MessageBuilder.withPayload(testPayload) + .setHeader(MqttHeaders.TOPIC, "anotherTopic") + .build()); + + receive = this.fromMqttChannel.receive(10_000); + + assertThat(receive).isNotNull(); + assertThat(receive.getPayload()).isEqualTo(testPayload); }