Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 {
Expand All @@ -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();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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) {
Expand Down Expand Up @@ -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);
}
}
Original file line number Diff line number Diff line change
@@ -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;
}
}
12 changes: 12 additions & 0 deletions docs/user-manual/broker-plugins.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -205,4 +205,16 @@ The plugin can be configured via xml in the normal broker-plugin way:
<property key="acceptorMatchRegex" value="netty-client-acceptor" />
</broker-plugin>
</broker-plugins>
----

== 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]
----
<broker-plugins>
<broker-plugin class-name="org.apache.activemq.artemis.core.protocol.mqtt.MQTTRetainMessagePlugin" />
</broker-plugins>
----
4 changes: 4 additions & 0 deletions docs/user-manual/mqtt.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading