From 57529cc38915ed8cbb9202770206d3635c3c357e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?V=C3=A1clav=20Haisman?= Date: Fri, 2 Feb 2024 15:00:11 -0500 Subject: [PATCH] GH-8879: Add MQTT subscription identifier Fixes: #8879 To work around the problem with `$share/` subscriptions the `Mqttv5PahoMessageDrivenChannelAdapter` must provide a `subscriptionIdentifier` into `MqttProperties` on `subscribe()` * Introduce a `Mqttv5PahoMessageDrivenChannelAdapter.subscriptionIdentifierCounter` according to the MQTT specification: > 3.8.2.1.2 Subscription Identifier: [..]The Subscription Identifier is associated with any subscription created or modified as the result of this SUBSCRIBE packet. If there is a Subscription Identifier, it is stored with the subscription. This one is associated with the MQTT session for the current subscriber and does not interfere into other sessions even if identifier is same from the counter. It works because the Subscription identifier is per session and because you cannot have multiple connection with the same client ID. **Cherry-pick to `6.2.x` & `6.1.x`** # Conflicts: # spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java --- .../Mqttv5PahoMessageDrivenChannelAdapter.java | 13 ++++++++++--- .../integration/mqtt/Mqttv5BackToBackTests.java | 16 ++++++++++++++++ 2 files changed, 26 insertions(+), 3 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 330b37aa4d..e35e4df7f2 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 @@ -1,5 +1,5 @@ /* - * Copyright 2021-2023 the original author or authors. + * Copyright 2021-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. @@ -18,6 +18,7 @@ package org.springframework.integration.mqtt.inbound; import java.util.Arrays; import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.stream.IntStream; @@ -100,7 +101,7 @@ public class Mqttv5PahoMessageDrivenChannelAdapter private volatile boolean readyToSubscribeOnStart; - public Mqttv5PahoMessageDrivenChannelAdapter(String url, String clientId, String... topic) { + private final AtomicInteger subscriptionIdentifierCounter = new AtomicInteger(0); public Mqttv5PahoMessageDrivenChannelAdapter(String url, String clientId, String... topic) { super(url, clientId, topic); Assert.hasText(url, "'url' cannot be null or empty"); this.connectionOptions = new MqttConnectionOptions(); @@ -276,8 +277,10 @@ public class Mqttv5PahoMessageDrivenChannelAdapter try { super.addTopic(topic, qos); if (this.mqttClient != null && this.mqttClient.isConnected()) { + MqttProperties subscriptionProperties = new MqttProperties(); + subscriptionProperties.setSubscriptionIdentifier(this.subscriptionIdentifierCounter.incrementAndGet()); this.mqttClient.subscribe(new MqttSubscription[] { new MqttSubscription(topic, qos) }, - null, null, new IMqttMessageListener[] { this::messageArrived }, new MqttProperties()) + null, null, new IMqttMessageListener[] { this::messageArrived }, subscriptionProperties) .waitForCompletion(getCompletionTimeout()); } } @@ -409,6 +412,10 @@ public class Mqttv5PahoMessageDrivenChannelAdapter IMqttMessageListener[] listeners = IntStream.range(0, topics.length) .mapToObj(t -> listener) .toArray(IMqttMessageListener[]::new); + MqttProperties subscriptionProperties = new MqttProperties(); + subscriptionProperties.setSubscriptionIdentifiers(IntStream.range(0, topics.length) + .mapToObj(i -> this.subscriptionIdentifierCounter.incrementAndGet()) + .toList()); this.mqttClient.subscribe(subscriptions, null, null, listeners, new MqttProperties()) .waitForCompletion(getCompletionTimeout()); String message = "Connected and subscribed to " + Arrays.toString(topics); 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 82f696d369..c9bf5d1f32 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 @@ -134,6 +134,22 @@ public class Mqttv5BackToBackTests implements MosquittoContainerTest { assertThat(receive.getPayload()).isEqualTo(testPayload); } + @Test + public void testSharedTopicMqttv5Interaction() { + this.mqttv5MessageDrivenChannelAdapter.addTopic("$share/group/testTopic"); + + String testPayload = "shared topic payload"; + this.mqttOutFlowInput.send( + MessageBuilder.withPayload(testPayload) + .setHeader(MqttHeaders.TOPIC, "testTopic") + .build()); + + Message receive = this.fromMqttChannel.receive(10_000); + + assertThat(receive).isNotNull(); + assertThat(receive.getPayload()).isEqualTo(testPayload); + } + @Configuration @EnableIntegration