diff --git a/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTPublishManager.java b/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTPublishManager.java index 5e6253e7efa..b35e559a5a5 100644 --- a/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTPublishManager.java +++ b/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTPublishManager.java @@ -78,7 +78,7 @@ public class MQTTPublishManager { private final MQTTSession session; - private boolean closeMqttConnectionOnPublishAuthorizationFailure; + private final boolean closeMqttConnectionOnPublishAuthorizationFailure; public MQTTPublishManager(MQTTSession session, boolean closeMqttConnectionOnPublishAuthorizationFailure) { this.session = session; @@ -164,7 +164,7 @@ synchronized void sendToQueue(MqttPublishMessage message, boolean internal) thro serverMessage.setDurable(MQTTUtil.DURABLE_MESSAGES); } - // only start a transction if really necessary + // only start a transaction if really necessary Transaction tx = realQos2 || message.fixedHeader().isRetain() ? session.getServerSession().newTransaction() : null; try { @@ -188,6 +188,7 @@ synchronized void sendToQueue(MqttPublishMessage message, boolean internal) thro boolean reset = payload instanceof EmptyByteBuf || payload.capacity() == 0; session.getRetainMessageManager().handleRetainedMessage(serverMessage, topic, reset, tx); } + if (tx != null) { tx.commit(); } diff --git a/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessageManager.java b/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessageManager.java index 6dc6d9ac290..29fc5410dc3 100644 --- a/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessageManager.java +++ b/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessageManager.java @@ -24,8 +24,6 @@ import org.apache.activemq.artemis.core.server.BindingQueryResult; import org.apache.activemq.artemis.core.server.MessageReference; import org.apache.activemq.artemis.core.server.Queue; -import org.apache.activemq.artemis.core.server.RoutingContext; -import org.apache.activemq.artemis.core.server.impl.RoutingContextImpl; import org.apache.activemq.artemis.core.transaction.Transaction; import org.apache.activemq.artemis.utils.collections.LinkedListIterator; @@ -34,28 +32,36 @@ public class MQTTRetainMessageManager { private MQTTSession session; + private final boolean retainMessagePluginRegistered; public MQTTRetainMessageManager(MQTTSession session) { this.session = session; + this.retainMessagePluginRegistered = session.getServer().getBrokerPlugins().stream().anyMatch(activeMQServerBasePlugin -> activeMQServerBasePlugin instanceof MQTTRetainMessagePlugin); } /** - * FIXME - * Retained messages should be handled in the core API. There is currently no support for retained messages - * at the time of writing. Instead we handle retained messages here. This method will create a new queue for - * every address that is used to store retained messages. THere should only ever be one message in the retained - * message queue. When a new subscription is created the queue should be browsed and the message copied onto + * Retained messages have two implementations, one that works on the mqtt protocol and is limited to a single broker + * and a second implemented as a broker plugin that intercepts all messages and can be used with broker connections. + * The tradeoff is that the plugin intercepts every message looking for the retain header, the plugin should only be + * configured when mqtt retained state needs to propagate between brokers. + * + * The implementation will create a new queue for every address that is used to store retained messages. + * There should only ever be one message in the retained message queue. + * When a new subscription is created the queue should be browsed and the message copied onto * the subscription queue for the consumer. When a new retained message is received the message will be sent to * the retained queue and the previous retain message consumed to remove it from the queue. */ void handleRetainedMessage(Message messageParameter, String address, boolean reset, Transaction tx) throws Exception { - String retainAddress = MQTTUtil.getCoreRetainAddressFromMqttTopic(address, session.getWildcardConfiguration()); - Queue queue = session.getServer().locateQueue(retainAddress); - if (queue == null) { - queue = session.getServer().createQueue(QueueConfiguration.of(retainAddress).setAutoCreated(true)); + if (retainMessagePluginRegistered) { + // see: org.apache.activemq.artemis.core.protocol.mqtt.MQTTRetainMessagePlugin.beforeMessageRoute + return; } + final String retainAddress = MQTTUtil.getCoreRetainAddressFromMqttTopic(address, session.getWildcardConfiguration()); + + Queue queue = session.getServer().createQueue(QueueConfiguration.of(retainAddress).setAutoCreated(true), true); + queue.deleteAllReferences(); if (!reset) { @@ -96,10 +102,4 @@ void addRetainedMessagesToQueue(Queue queue, String address) throws Exception { } tx.commit(); } - - private void sendToQueue(Message message, Queue queue, Transaction tx) throws Exception { - RoutingContext context = new RoutingContextImpl(tx); - queue.route(message, context); - session.getServer().getPostOffice().processRoute(message, context, false); - } } diff --git a/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessagePlugin.java b/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessagePlugin.java new file mode 100644 index 00000000000..935f17ae444 --- /dev/null +++ b/artemis-protocols/artemis-mqtt-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/mqtt/MQTTRetainMessagePlugin.java @@ -0,0 +1,85 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 + * + * http://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.apache.activemq.artemis.core.protocol.mqtt; + +import org.apache.activemq.artemis.api.core.ActiveMQException; +import org.apache.activemq.artemis.api.core.Message; +import org.apache.activemq.artemis.api.core.QueueConfiguration; +import org.apache.activemq.artemis.core.persistence.StorageManager; +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.server.Queue; +import org.apache.activemq.artemis.core.server.RoutingContext; +import org.apache.activemq.artemis.core.server.plugin.ActiveMQServerMessagePlugin; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import java.lang.invoke.MethodHandles; + +public class MQTTRetainMessagePlugin implements ActiveMQServerMessagePlugin { + + private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); + + private ActiveMQServer server; + + @Override + public void registered(ActiveMQServer server) { + this.server = server; + } + + /** + * reacting to a retain header on all messages published, useful when a mqtt message is forwarded from + * another broker. It won't use the mqtt protocol head in that case. + * The MQTTRetainMessageManager delegates to this plugin when it is present in broker config. + * see org.apache.activemq.artemis.core.protocol.mqtt.MQTTRetainMessageManager#handleRetainedMessage + */ + @Override + public void beforeMessageRoute(Message message, RoutingContext context, boolean direct, boolean rejectDuplicates) throws ActiveMQException { + try { + Boolean isRetain = message.getBooleanProperty(MQTTUtil.MQTT_MESSAGE_RETAIN_KEY); + if (isRetain == null || !isRetain) { + return; + } + + final String address = message.getAddress(); + if (address == null) { + return; + } + + final String retainAddress = MQTTUtil.MQTT_RETAIN_ADDRESS_PREFIX + address; + final Queue queue = server.createQueue(QueueConfiguration.of(retainAddress).setAutoCreated(true), true); + + queue.deleteAllReferences(); + + // retain only if non-empty + if (message.isLargeMessage() || message.toCore().getBodyBufferSize() > 0) { + final StorageManager storageManager = server.getStorageManager(); + MQTTUtil.sendMessageDirectlyToQueue(storageManager, server.getPostOffice(), message.copy(storageManager.generateID()), queue, context.getTransaction()); + } + } catch (Exception e) { + logger.warn("Failed to handle MQTT retained message for address {}: {}", message.getAddress(), e.getMessage(), e); + } + } + + @Override + public int hashCode() { + return System.identityHashCode(MQTTRetainMessagePlugin.class); + } + + @Override + public boolean equals(Object obj) { + return obj instanceof MQTTRetainMessagePlugin; + } +} diff --git a/docs/user-manual/broker-plugins.adoc b/docs/user-manual/broker-plugins.adoc index 3345b7032c5..472220306a7 100644 --- a/docs/user-manual/broker-plugins.adoc +++ b/docs/user-manual/broker-plugins.adoc @@ -205,4 +205,16 @@ The plugin can be configured via xml in the normal broker-plugin way: +---- + +== Using the MQTTRetainMessagePlugin + +The `MQTTRetainMessagePlugin` will move xref:mqtt.adoc#mqtt-retain-messages[MQTT retain message processing] from the MQTT publish manager on the MQTT protocol head to the broker routing layer. Every message will be checked for the MQTT specific retain header such that MQTT messages that orginate on the broker over other protocols (via an AMQPBrokerConnection for example) will have their retain header respected. The plugin checks the properties of every message, so this plugin should only be used when the MQTT retain state needs to span more than one broker. + +The plugin can be configured via xml in the normal broker-plugin way: +[,xml] +---- + + + ---- \ No newline at end of file diff --git a/docs/user-manual/mqtt.adoc b/docs/user-manual/mqtt.adoc index 45530c5c783..71ae3edc7b3 100644 --- a/docs/user-manual/mqtt.adoc +++ b/docs/user-manual/mqtt.adoc @@ -72,6 +72,10 @@ These resources can be automatically deleted via the following `address-setting` Keep in mind that it's also possible to automatically apply an xref:message-expiry.adoc#message-expiry[`expiry-delay`] to retained messages as well. +=== Retain messages with AMQP xref:amqp-broker-connections.adoc#address-federation[address federation] +If retain messages need to work across a group of brokers and those brokers are linked with an xref:amqp-broker-connections.adoc#federation[AMQPBrokerConnection], it is necessary to move the retain processing from the MQTT protocol head into the protocol independent routing layer of the Broker. +There is a xref:broker-plugins.adoc#using-the-mqttretainmessageplugin[MQTTRetainMessagePlugin] for that. One example use case would be partitioned pub/sub, where consumers are partitioned across brokers and publishers can use any broker. To achieve the desired publishing fanout, all addresses are federated. The receiving end of the amqp federation will need to respect any retain header for its broker to store retained messages. + == Will Messages A will message can be sent when a client initially connects to a broker. diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5BridgeTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5BridgeTest.java new file mode 100644 index 00000000000..11f7a0b6b62 --- /dev/null +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5BridgeTest.java @@ -0,0 +1,215 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 + * + * http://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.apache.activemq.artemis.tests.integration.mqtt5; + +import java.nio.charset.StandardCharsets; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.activemq.artemis.core.protocol.mqtt.MQTTRetainMessagePlugin; +import org.apache.activemq.artemis.tests.util.ActiveMQTestBase; +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBridgeAddressPolicyElement; +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBridgeBrokerConnectionElement; +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBrokerConnectConfiguration; +import org.apache.activemq.artemis.core.protocol.mqtt.MQTTUtil; +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.settings.impl.AddressSettings; +import org.apache.activemq.artemis.utils.Wait; +import org.eclipse.paho.mqttv5.client.MqttClient; +import org.eclipse.paho.mqttv5.client.persist.MemoryPersistence; +import org.eclipse.paho.mqttv5.common.MqttException; +import org.eclipse.paho.mqttv5.common.MqttMessage; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class MQTT5BridgeTest extends ActiveMQTestBase { + + private static final int BROKER1_PORT = 1883; + private static final int BROKER2_PORT = 1884; + + @BeforeEach + @Override + public void setUp() throws Exception { + super.setUp(); + } + + private ActiveMQServer createServer(int serverID, int port) throws Exception { + ActiveMQServer server = createServer(true, createDefaultConfig(serverID, true)); + server.getConfiguration().getAcceptorConfigurations().clear(); + server.getConfiguration().addAcceptorConfiguration("server", "tcp://localhost:" + port); + server.getConfiguration().setSecurityEnabled(false); + server.getConfiguration().setMqttSessionScanInterval(200); + server.getConfiguration().registerBrokerPlugin(new MQTTRetainMessagePlugin()); + + AddressSettings addressSettings = new AddressSettings(); + addressSettings.setAutoCreateQueues(true); + addressSettings.setAutoCreateAddresses(true); + server.getConfiguration().getAddressSettings().put("#", addressSettings); + + return server; + } + + private ActiveMQServer createBridgedServer() throws Exception { + ActiveMQServer server = createServer(1, MQTT5BridgeTest.BROKER2_PORT); + + final AMQPBridgeAddressPolicyElement addressPolicy = new AMQPBridgeAddressPolicyElement(); + addressPolicy.setName("mqtt-bridge-policy"); + addressPolicy.addToIncludes("#"); + + final AMQPBridgeBrokerConnectionElement bridgeElement = new AMQPBridgeBrokerConnectionElement(); + bridgeElement.setName("bridgeFromServer1"); + bridgeElement.addBridgeFromAddressPolicy(addressPolicy); + + final AMQPBrokerConnectConfiguration amqpConnection = + new AMQPBrokerConnectConfiguration("bridgeFromServer1", "tcp://localhost:" + MQTT5BridgeTest.BROKER1_PORT) + .setReconnectAttempts(10) + .setRetryInterval(100); + amqpConnection.addElement(bridgeElement); + + server.getConfiguration().addAMQPConnection(amqpConnection); + + return server; + } + + private MqttClient createPahoClient(String clientId, int port) throws MqttException { + return new MqttClient("tcp://localhost:" + port, clientId, new MemoryPersistence()); + } + + @Test + @Timeout(60) + public void testRetainedMessageOverBridge() throws Exception { + // broker 1 is a plain server (source); broker 2 bridges FROM broker 1 + ActiveMQServer server1 = createServer(0, BROKER1_PORT); + ActiveMQServer server2 = createBridgedServer(); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + final String topic = "test/retain/bridge"; + final String payload = "retained-via-bridge"; + + // subscribe on broker 2 first — this creates demand so the bridge pulls from broker 1 + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + + // publish retained message on broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + producer.disconnect(); + producer.close(); + + // subscriber on broker 2 should receive the message via the bridge + assertTrue(latch.await(10, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the bridged message"); + assertEquals(payload, received.get()); + + // verify retain queue on server1 (local MQTT retain handling) + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + // retain queue should also exist on server2 — populated by the plugin + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + sub.disconnect(); + sub.close(); + + // a new subscriber on broker 2 should get the retained message from the local retain queue + CountDownLatch retainLatch = new CountDownLatch(1); + AtomicReference retainReceived = new AtomicReference<>(); + MqttClient sub2 = createPahoClient("subscriber2", BROKER2_PORT); + sub2.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + retainReceived.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + retainLatch.countDown(); + } + }); + sub2.connect(); + sub2.subscribe(topic, 1); + assertTrue(retainLatch.await(5, TimeUnit.SECONDS), "New subscriber on broker 2 should receive the retained message"); + assertEquals(payload, retainReceived.get()); + + sub2.disconnect(); + sub2.close(); + } + + @Test + @Timeout(60) + public void testRegularMessageOverBridge() throws Exception { + ActiveMQServer server1 = createServer(0, BROKER1_PORT); + ActiveMQServer server2 = createBridgedServer(); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + final String topic = "test/bridge/regular"; + final String payload = "regular-via-bridge"; + + // subscribe on broker 2 to create demand + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + + // publish a regular (non-retained) message on broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, false); + producer.disconnect(); + producer.close(); + + // subscriber on broker 2 should receive the message via the bridge + assertTrue(latch.await(10, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the bridged message"); + assertEquals(payload, received.get()); + + sub.disconnect(); + sub.close(); + } +} diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5FederationTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5FederationTest.java new file mode 100644 index 00000000000..457ca387b19 --- /dev/null +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5FederationTest.java @@ -0,0 +1,353 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 + * + * http://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.apache.activemq.artemis.tests.integration.mqtt5; + +import java.nio.charset.StandardCharsets; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBrokerConnectConfiguration; +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPFederatedBrokerConnectionElement; +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPFederationAddressPolicyElement; +import org.apache.activemq.artemis.core.protocol.mqtt.MQTTRetainMessagePlugin; +import org.apache.activemq.artemis.core.protocol.mqtt.MQTTUtil; +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.settings.impl.AddressSettings; +import org.apache.activemq.artemis.tests.util.ActiveMQTestBase; +import org.apache.activemq.artemis.utils.Wait; +import org.eclipse.paho.mqttv5.client.MqttClient; +import org.eclipse.paho.mqttv5.client.persist.MemoryPersistence; +import org.eclipse.paho.mqttv5.common.MqttException; +import org.eclipse.paho.mqttv5.common.MqttMessage; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class MQTT5FederationTest extends ActiveMQTestBase { + + private static final int BROKER1_PORT = 1883; + private static final int BROKER2_PORT = 1884; + + @BeforeEach + @Override + public void setUp() throws Exception { + super.setUp(); + } + + private ActiveMQServer createFederatedServer(int serverID, int port, int targetPort, String connectionName) throws Exception { + ActiveMQServer server = createServer(true, createDefaultConfig(serverID, true)); + server.getConfiguration().registerBrokerPlugin(new MQTTRetainMessagePlugin()); + server.getConfiguration().getAcceptorConfigurations().clear(); + server.getConfiguration().addAcceptorConfiguration("server", "tcp://localhost:" + port); + server.getConfiguration().setSecurityEnabled(false); + server.getConfiguration().setMqttSessionScanInterval(200); + + AddressSettings addressSettings = new AddressSettings(); + addressSettings.setAutoCreateQueues(true); + addressSettings.setAutoCreateAddresses(true); + server.getConfiguration().getAddressSettings().put("#", addressSettings); + + final AMQPFederationAddressPolicyElement localAddressPolicy = new AMQPFederationAddressPolicyElement(); + localAddressPolicy.setName("mqtt-address-policy"); + localAddressPolicy.addToIncludes("#"); + localAddressPolicy.setAutoDelete(false); + localAddressPolicy.setAutoDeleteDelay(-1L); + localAddressPolicy.setAutoDeleteMessageCount(-1L); + + final AMQPFederatedBrokerConnectionElement federationElement = new AMQPFederatedBrokerConnectionElement(); + federationElement.setName(connectionName); + federationElement.addLocalAddressPolicy(localAddressPolicy); + + final AMQPBrokerConnectConfiguration amqpConnection = + new AMQPBrokerConnectConfiguration(connectionName, "tcp://localhost:" + targetPort) + .setReconnectAttempts(10) + .setRetryInterval(100); + amqpConnection.addElement(federationElement); + + server.getConfiguration().addAMQPConnection(amqpConnection); + + return server; + } + + + private MqttClient createPahoClient(String clientId, int port) throws MqttException { + return new MqttClient("tcp://localhost:" + port, clientId, new MemoryPersistence()); + } + + @Test + @Timeout(60) + public void testPartitionedSubscriptionsWithFederation() throws Exception { + ActiveMQServer server1 = createFederatedServer(0, BROKER1_PORT, BROKER2_PORT, "federationToServer2"); + ActiveMQServer server2 = createFederatedServer(1, BROKER2_PORT, BROKER1_PORT, "federationToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + final String sportsTopic = "news/sports"; + final String techTopic = "news/tech"; + + // subscribers are partitioned: sub-a on broker 1 for sports, sub-b on broker 2 for tech + CountDownLatch sportsLatch = new CountDownLatch(2); + CopyOnWriteArrayList sportsReceived = new CopyOnWriteArrayList<>(); + MqttClient subA = createPahoClient("sub-a", BROKER1_PORT); + subA.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + sportsReceived.add(new String(m.getPayload(), StandardCharsets.UTF_8)); + sportsLatch.countDown(); + } + }); + subA.connect(); + subA.subscribe(sportsTopic, 1); + + CountDownLatch techLatch = new CountDownLatch(2); + CopyOnWriteArrayList techReceived = new CopyOnWriteArrayList<>(); + MqttClient subB = createPahoClient("sub-b", BROKER2_PORT); + subB.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + techReceived.add(new String(m.getPayload(), StandardCharsets.UTF_8)); + techLatch.countDown(); + } + }); + subB.connect(); + subB.subscribe(techTopic, 1); + + // publish on the broker that has NO local subscriber for that topic + MqttClient pub1 = createPahoClient("pub1", BROKER2_PORT); + pub1.connect(); + pub1.publish(sportsTopic, "goal".getBytes(StandardCharsets.UTF_8), 1, false); + pub1.disconnect(); + pub1.close(); + + MqttClient pub2 = createPahoClient("pub2", BROKER1_PORT); + pub2.connect(); + pub2.publish(techTopic, "update".getBytes(StandardCharsets.UTF_8), 1, false); + pub2.disconnect(); + pub2.close(); + + // also publish locally to verify local delivery still works + MqttClient pub3 = createPahoClient("pub3", BROKER1_PORT); + pub3.connect(); + pub3.publish(sportsTopic, "try".getBytes(StandardCharsets.UTF_8), 1, false); + pub3.disconnect(); + pub3.close(); + + MqttClient pub4 = createPahoClient("pub4", BROKER2_PORT); + pub4.connect(); + pub4.publish(techTopic, "release".getBytes(StandardCharsets.UTF_8), 1, false); + pub4.disconnect(); + pub4.close(); + + assertTrue(sportsLatch.await(10, TimeUnit.SECONDS), "sub-a should receive both sports messages"); + assertTrue(techLatch.await(10, TimeUnit.SECONDS), "sub-b should receive both tech messages"); + + assertEquals(2, sportsReceived.size()); + assertTrue(sportsReceived.contains("goal"), "sub-a should receive 'goal' published on broker 2"); + assertTrue(sportsReceived.contains("try"), "sub-a should receive 'try' published on broker 1"); + + assertEquals(2, techReceived.size()); + assertTrue(techReceived.contains("update"), "sub-b should receive 'update' published on broker 1"); + assertTrue(techReceived.contains("release"), "sub-b should receive 'release' published on broker 2"); + + subA.disconnect(); + subA.close(); + subB.disconnect(); + subB.close(); + } + + @Test + @Timeout(60) + public void testRetainedMessageWithFederation() throws Exception { + ActiveMQServer server1 = createFederatedServer(0, BROKER1_PORT, BROKER2_PORT, "federationToServer2"); + ActiveMQServer server2 = createFederatedServer(1, BROKER2_PORT, BROKER1_PORT, "federationToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + final String topic = "test/retain/federation"; + final String payload = "retained-message-payload"; + + // subscribe on broker 2 first — this creates federation demand + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + + // publish retained message on broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + producer.disconnect(); + producer.close(); + + // subscriber on broker 2 should receive the message via federation + assertTrue(latch.await(10, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the federated message"); + assertEquals(payload, received.get()); + + // verify retain queue exists on server1 (local retain handling) + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + // retain queue should also exist on server2 — populated by federation consumer + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + // verify no duplicate retain processing: messagesAdded should be exactly 1 + assertEquals(1, server2.locateQueue(retainQueueName).getMessagesAdded(), + "Retain queue on server2 should have exactly 1 message added (no duplicate processing)"); + + sub.disconnect(); + sub.close(); + + // a new subscriber on broker 2 should get the retained message from the local retain queue + CountDownLatch retainLatch = new CountDownLatch(1); + AtomicReference retainReceived = new AtomicReference<>(); + MqttClient sub2 = createPahoClient("subscriber2", BROKER2_PORT); + sub2.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + retainReceived.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + retainLatch.countDown(); + } + }); + sub2.connect(); + sub2.subscribe(topic, 1); + assertTrue(retainLatch.await(5, TimeUnit.SECONDS), "New subscriber on broker 2 should receive the retained message"); + assertEquals(payload, retainReceived.get()); + + sub2.disconnect(); + sub2.close(); + } + + @Test + @Timeout(60) + public void testClearRetainedMessageWithFederation() throws Exception { + ActiveMQServer server1 = createFederatedServer(0, BROKER1_PORT, BROKER2_PORT, "federationToServer2"); + ActiveMQServer server2 = createFederatedServer(1, BROKER2_PORT, BROKER1_PORT, "federationToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + final String topic = "test/retain/federation/clear"; + final String payload = "retained-to-clear"; + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + + // subscribe on broker 2 to create federation demand + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + + // publish a retained message on broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + + assertTrue(latch.await(10, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the federated message"); + assertEquals(payload, received.get()); + + // verify retain queue populated on both brokers + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + // clear the retained message with a zero-length payload (MQTT spec §3.3.1.3) + producer.publish(topic, new byte[0], 1, true); + producer.disconnect(); + producer.close(); + + // retain queue on server1 should be empty + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 0; + }, 5000, 100); + + // retain queue on server2 should also be empty — cleared via federated zero-length message + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 0; + }, 10000, 100); + + sub.disconnect(); + sub.close(); + + // a new subscriber on broker 2 should NOT receive a retained message + CountDownLatch noRetainLatch = new CountDownLatch(1); + AtomicReference noRetainReceived = new AtomicReference<>(); + MqttClient sub2 = createPahoClient("subscriber2", BROKER2_PORT); + sub2.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + noRetainReceived.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + noRetainLatch.countDown(); + } + }); + sub2.connect(); + sub2.subscribe(topic, 1); + + // publish a regular message so we know the subscription is active + MqttClient probe = createPahoClient("probe", BROKER1_PORT); + probe.connect(); + probe.publish(topic, "probe".getBytes(StandardCharsets.UTF_8), 1, false); + probe.disconnect(); + probe.close(); + + assertTrue(noRetainLatch.await(10, TimeUnit.SECONDS), "Subscriber should receive the probe message"); + assertEquals("probe", noRetainReceived.get(), "Only the probe message should arrive — no retained message"); + + sub2.disconnect(); + sub2.close(); + } +} diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5MirrorTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5MirrorTest.java new file mode 100644 index 00000000000..8e2e9706fd5 --- /dev/null +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5MirrorTest.java @@ -0,0 +1,400 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 + * + * http://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.apache.activemq.artemis.tests.integration.mqtt5; + +import java.nio.charset.StandardCharsets; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBrokerConnectConfiguration; +import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPMirrorBrokerConnectionElement; +import org.apache.activemq.artemis.core.protocol.mqtt.MQTTRetainMessagePlugin; +import org.apache.activemq.artemis.core.protocol.mqtt.MQTTUtil; +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.settings.impl.AddressSettings; +import org.apache.activemq.artemis.tests.util.ActiveMQTestBase; +import org.apache.activemq.artemis.utils.Wait; +import org.eclipse.paho.mqttv5.client.MqttClient; +import org.eclipse.paho.mqttv5.client.persist.MemoryPersistence; +import org.eclipse.paho.mqttv5.common.MqttException; +import org.eclipse.paho.mqttv5.common.MqttMessage; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class MQTT5MirrorTest extends ActiveMQTestBase { + + private static final int BROKER1_PORT = 1883; + private static final int BROKER2_PORT = 1884; + + @BeforeEach + @Override + public void setUp() throws Exception { + super.setUp(); + } + + private ActiveMQServer createMirroredServer(int serverID, int port, int targetPort, String connectionName) throws Exception { + ActiveMQServer server = createServer(true, createDefaultConfig(serverID, true)); + server.getConfiguration().registerBrokerPlugin(new MQTTRetainMessagePlugin()); + server.getConfiguration().getAcceptorConfigurations().clear(); + server.getConfiguration().addAcceptorConfiguration("server", "tcp://localhost:" + port); + server.getConfiguration().setSecurityEnabled(false); + server.getConfiguration().setMqttSessionScanInterval(200); + + AddressSettings addressSettings = new AddressSettings(); + addressSettings.setAutoCreateQueues(true); + addressSettings.setAutoCreateAddresses(true); + server.getConfiguration().getAddressSettings().put("#", addressSettings); + + AMQPBrokerConnectConfiguration amqpConnection = new AMQPBrokerConnectConfiguration(connectionName, "tcp://localhost:" + targetPort) + .setReconnectAttempts(-1) + .setRetryInterval(100); + amqpConnection.addElement(new AMQPMirrorBrokerConnectionElement().setDurable(true)); + server.getConfiguration().addAMQPConnection(amqpConnection); + + return server; + } + + private MqttClient createPahoClient(String clientId, int port) throws MqttException { + return new MqttClient("tcp://localhost:" + port, clientId, new MemoryPersistence()); + } + + @Test + @Timeout(60) + public void testRetainedMessageMirrored() throws Exception { + ActiveMQServer server1 = createMirroredServer(0, BROKER1_PORT, BROKER2_PORT, "mirrorToServer2"); + ActiveMQServer server2 = createMirroredServer(1, BROKER2_PORT, BROKER1_PORT, "mirrorToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + Wait.assertTrue(() -> server1.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer2") != null); + Wait.assertTrue(() -> server2.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer1") != null); + + final String topic = "test/retain/mirror"; + final String payload = "retained-message-payload"; + + // publish retained message to broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + producer.disconnect(); + producer.close(); + + // verify the retain queue exists on server1 + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + // subscribe on broker 1 - should get the retained message + CountDownLatch latch1 = new CountDownLatch(1); + AtomicReference received1 = new AtomicReference<>(); + MqttClient sub1 = createPahoClient("subscriber1", BROKER1_PORT); + sub1.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received1.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch1.countDown(); + } + }); + sub1.connect(); + sub1.subscribe(topic, 1); + assertTrue(latch1.await(5, TimeUnit.SECONDS), "Subscriber on broker 1 should receive the retained message"); + assertEquals(payload, received1.get()); + sub1.disconnect(); + sub1.close(); + + // wait for the retained message to appear on broker 2 via mirroring + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 10000, 100); + + // subscribe on broker 2 - should also get the retained message + CountDownLatch latch2 = new CountDownLatch(1); + AtomicReference received2 = new AtomicReference<>(); + MqttClient sub2 = createPahoClient("subscriber2", BROKER2_PORT); + sub2.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received2.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch2.countDown(); + } + }); + sub2.connect(); + sub2.subscribe(topic, 1); + assertTrue(latch2.await(5, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the mirrored retained message"); + assertEquals(payload, received2.get()); + sub2.disconnect(); + sub2.close(); + + // publish a second retained message - should replace the first on both brokers + final String payload2 = "retained-message-payload-2"; + MqttClient producer2 = createPahoClient("producer2", BROKER1_PORT); + producer2.connect(); + producer2.publish(topic, payload2.getBytes(StandardCharsets.UTF_8), 1, true); + producer2.disconnect(); + producer2.close(); + + // verify server1 retain queue replaced: still exactly 1 message with the new payload + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1 && queue.getMessagesAdded() == 2; + }, 5000, 100); + + // verify server2 retain queue replaced via mirroring: exactly 1 message + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 10000, 100); + + // subscribe on broker 1 - should get the second retained message + CountDownLatch latch3 = new CountDownLatch(1); + AtomicReference received3 = new AtomicReference<>(); + MqttClient sub3 = createPahoClient("subscriber3", BROKER1_PORT); + sub3.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received3.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch3.countDown(); + } + }); + sub3.connect(); + sub3.subscribe(topic, 1); + assertTrue(latch3.await(5, TimeUnit.SECONDS), "Subscriber on broker 1 should receive the second retained message"); + assertEquals(payload2, received3.get()); + sub3.disconnect(); + sub3.close(); + + // subscribe on broker 2 - should get the second retained message + CountDownLatch latch4 = new CountDownLatch(1); + AtomicReference received4 = new AtomicReference<>(); + MqttClient sub4 = createPahoClient("subscriber4", BROKER2_PORT); + sub4.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received4.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch4.countDown(); + } + }); + sub4.connect(); + sub4.subscribe(topic, 1); + assertTrue(latch4.await(5, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the second retained message"); + assertEquals(payload2, received4.get()); + sub4.disconnect(); + sub4.close(); + } + + @Test + @Timeout(60) + public void testClearRetainedMessageMirrored() throws Exception { + ActiveMQServer server1 = createMirroredServer(0, BROKER1_PORT, BROKER2_PORT, "mirrorToServer2"); + ActiveMQServer server2 = createMirroredServer(1, BROKER2_PORT, BROKER1_PORT, "mirrorToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + Wait.assertTrue(() -> server1.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer2") != null); + Wait.assertTrue(() -> server2.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer1") != null); + + final String topic = "test/retain/clear"; + final String payload = "retained-to-clear"; + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + + // publish a retained message on broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + + // verify retain queue populated on both brokers + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 10000, 100); + + // publish a zero-length retained message to clear the retain (MQTT spec §3.3.1.3) + producer.publish(topic, new byte[0], 1, true); + producer.disconnect(); + producer.close(); + + // retain queue on server1 should be empty (cleared locally by MQTTRetainMessageManager) + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 0; + }, 5000, 100); + + // retain queue on server2 should also be empty (cleared via mirrored zero-length message) + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 0; + }, 10000, 100); + + // a new subscriber on broker 2 should NOT receive a retained message + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + + // publish a regular (non-retained) message so we know the subscription is active + MqttClient probe = createPahoClient("probe", BROKER2_PORT); + probe.connect(); + probe.publish(topic, "probe".getBytes(StandardCharsets.UTF_8), 1, false); + probe.disconnect(); + probe.close(); + + assertTrue(latch.await(5, TimeUnit.SECONDS), "Subscriber should receive the probe message"); + assertEquals("probe", received.get(), "Only the probe message should arrive — no retained message"); + + sub.disconnect(); + sub.close(); + } + + @Test + @Timeout(60) + public void testRetainedMessageSurvivesBrokerStop() throws Exception { + ActiveMQServer server1 = createMirroredServer(0, BROKER1_PORT, BROKER2_PORT, "mirrorToServer2"); + ActiveMQServer server2 = createMirroredServer(1, BROKER2_PORT, BROKER1_PORT, "mirrorToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + Wait.assertTrue(() -> server1.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer2") != null); + Wait.assertTrue(() -> server2.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer1") != null); + + final String topic = "test/retain/failover"; + final String payload = "retained-survives-stop"; + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + + // publish retained on broker 1 + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + producer.disconnect(); + producer.close(); + + // wait for retain queue on both brokers + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 10000, 100); + + // stop broker 1 + server1.stop(); + + // a new subscriber on broker 2 should still get the retained message from the local retain queue + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + assertTrue(latch.await(5, TimeUnit.SECONDS), "Subscriber on broker 2 should receive the retained message after broker 1 stopped"); + assertEquals(payload, received.get()); + sub.disconnect(); + sub.close(); + } + + @Test + @Timeout(60) + public void testRetainedMessageMirroredWithNoLocalSubscribers() throws Exception { + ActiveMQServer server1 = createMirroredServer(0, BROKER1_PORT, BROKER2_PORT, "mirrorToServer2"); + ActiveMQServer server2 = createMirroredServer(1, BROKER2_PORT, BROKER1_PORT, "mirrorToServer1"); + + server1.start(); + server1.waitForActivation(10, TimeUnit.SECONDS); + server2.start(); + server2.waitForActivation(10, TimeUnit.SECONDS); + + Wait.assertTrue(() -> server1.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer2") != null); + Wait.assertTrue(() -> server2.locateQueue("$ACTIVEMQ_ARTEMIS_MIRROR_mirrorToServer1") != null); + + final String topic = "test/retain/nosub"; + final String payload = "retained-no-subscribers"; + final String retainQueueName = MQTTUtil.getCoreRetainAddressFromMqttTopic(topic, server1.getConfiguration().getWildcardConfiguration()); + + // publish a retained message on broker 1 with NO subscribers anywhere + MqttClient producer = createPahoClient("producer", BROKER1_PORT); + producer.connect(); + producer.publish(topic, payload.getBytes(StandardCharsets.UTF_8), 1, true); + producer.disconnect(); + producer.close(); + + // retain queue on server1 should be populated locally by the plugin + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server1.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 5000, 100); + + // retain queue on server2 should also be populated via create and mirror + Wait.assertTrue(() -> { + org.apache.activemq.artemis.core.server.Queue queue = server2.locateQueue(retainQueueName); + return queue != null && queue.getMessageCount() == 1; + }, 10000, 100); + + // a new subscriber on broker 2 should receive the retained message + CountDownLatch latch = new CountDownLatch(1); + AtomicReference received = new AtomicReference<>(); + MqttClient sub = createPahoClient("subscriber", BROKER2_PORT); + sub.setCallback(new MQTT5TestSupport.DefaultMqttCallback() { + @Override + public void messageArrived(String t, MqttMessage m) { + received.set(new String(m.getPayload(), StandardCharsets.UTF_8)); + latch.countDown(); + } + }); + sub.connect(); + sub.subscribe(topic, 1); + assertTrue(latch.await(5, TimeUnit.SECONDS), "New subscriber on broker 2 should receive the retained message"); + assertEquals(payload, received.get()); + sub.disconnect(); + sub.close(); + } +}