diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/ConsumeMQTT.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/ConsumeMQTT.java index 7763eb8a4c9c..74a420f64ec1 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/ConsumeMQTT.java +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/ConsumeMQTT.java @@ -45,6 +45,7 @@ import org.apache.nifi.processor.util.StandardValidators; import org.apache.nifi.processors.mqtt.common.AbstractMQTTProcessor; import org.apache.nifi.processors.mqtt.common.MqttException; +import org.apache.nifi.processors.mqtt.common.MqttTopicSubscription; import org.apache.nifi.processors.mqtt.common.ReceivedMqttMessage; import org.apache.nifi.serialization.MalformedRecordException; import org.apache.nifi.serialization.RecordReader; @@ -67,6 +68,8 @@ import java.util.ArrayList; import java.util.Collection; import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -75,6 +78,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; import static org.apache.nifi.processors.mqtt.ConsumeMQTT.BROKER_ATTRIBUTE_KEY; import static org.apache.nifi.processors.mqtt.ConsumeMQTT.IS_DUPLICATE_ATTRIBUTE_KEY; @@ -133,7 +137,10 @@ public class ConsumeMQTT extends AbstractMQTTProcessor { public static final PropertyDescriptor PROP_TOPIC_FILTER = new PropertyDescriptor.Builder() .name("Topic Filter") - .description("The MQTT topic filter to designate the topics to subscribe to.") + .description("The MQTT topic filter to designate the topics to subscribe to. More than one can be supplied if comma separated, in which case a single SUBSCRIBE request " + + "listing every filter is sent to the broker, avoiding the need for a separate processor and broker connection per topic. A value without a comma is used as a " + + "single topic filter exactly as configured, while the entries of a comma separated value are trimmed. Because MQTT topic filters may legally contain a comma, " + + "a Topic Filter containing one is interpreted as multiple filters; this is a rare edge case but should be kept in mind.") .required(true) .expressionLanguageSupported(ExpressionLanguageScope.ENVIRONMENT) .addValidator(StandardValidators.NON_BLANK_VALIDATOR) @@ -192,6 +199,7 @@ public class ConsumeMQTT extends AbstractMQTTProcessor { private volatile int qos; private volatile String topicPrefix = ""; private volatile String topicFilter; + private volatile List topicSubscriptions = List.of(); private final AtomicBoolean scheduled = new AtomicBoolean(false); private volatile BlockingQueue mqttQueue; @@ -288,6 +296,27 @@ public Collection customValidate(ValidationContext context) { .build()); } + final String rawTopicFilter = context.getProperty(PROP_TOPIC_FILTER).evaluateAttributeExpressions().getValue(); + if (rawTopicFilter != null) { + final List topicFilters = parseTopicFilters(rawTopicFilter); + if (topicFilters.isEmpty()) { + results.add(new ValidationResult.Builder() + .subject(PROP_TOPIC_FILTER.getDisplayName()) + .valid(false) + .explanation("at least one non-blank Topic Filter must be provided.") + .build()); + } else { + final Set uniqueTopicFilters = new HashSet<>(topicFilters); + if (uniqueTopicFilters.size() != topicFilters.size()) { + results.add(new ValidationResult.Builder() + .subject(PROP_TOPIC_FILTER.getDisplayName()) + .valid(false) + .explanation("duplicate Topic Filters are not allowed: " + topicFilters) + .build()); + } + } + } + return results; } @@ -319,6 +348,12 @@ public void onScheduled(final ProcessContext context) { topicPrefix = ""; } + // The shared subscription prefix applies to an individual Topic Filter, so it has to be added to each of them + // separately rather than to the configured, potentially comma separated, value as a whole. + topicSubscriptions = new LinkedHashSet<>(parseTopicFilters(topicFilter)).stream() + .map(filter -> new MqttTopicSubscription(topicPrefix + filter, qos)) + .collect(Collectors.toList()); + scheduled.set(true); } @@ -389,14 +424,39 @@ private void initializeClient(ProcessContext context) { try { mqttClient = createMqttClient(); mqttClient.connect(); - mqttClient.subscribe(topicPrefix + topicFilter, qos, this::handleReceivedMessage); + mqttClient.subscribe(topicSubscriptions, this::handleReceivedMessage); } catch (Exception e) { logger.error("Connection failed to {}. Yielding processor", clientProperties.getRawBrokerUris(), e); - mqttClient = null; // prevent stuck processor when subscribe fails + // A SUBSCRIBE carrying several Topic Filters can be granted partially, so the client may be connected and + // subscribed even though subscribe() failed. Disconnecting and closing it, rather than only dropping the + // reference, prevents an orphaned client from holding a broker connection and feeding the internal queue. + stopClient(); context.yield(); } } + /** + * Splits the configured Topic Filter property, which may contain a comma-separated list of topic filters, into + * a list of non-blank topic filters, preserving duplicates so that {@link #customValidate(ValidationContext)} can + * flag them. A value without a comma is a single topic filter and is used verbatim, so existing configurations + * behave identically. Only the segments of a comma-separated value are trimmed, because leading and trailing + * whitespace is significant in an MQTT topic filter and trimming is merely a convenience for writing a list. + */ + private static List parseTopicFilters(final String rawTopicFilters) { + if (rawTopicFilters.indexOf(',') < 0) { + return List.of(rawTopicFilters); + } + + final List topicFilters = new ArrayList<>(); + for (final String topicFilter : rawTopicFilters.split(",", -1)) { + final String trimmedTopicFilter = topicFilter.trim(); + if (!trimmedTopicFilter.isEmpty()) { + topicFilters.add(trimmedTopicFilter); + } + } + return topicFilters; + } + private void transferQueue(ProcessSession session) { while (!mqttQueue.isEmpty()) { final ReceivedMqttMessage mqttMessage = mqttQueue.peek(); @@ -430,7 +490,7 @@ private void transferQueueDemarcator(final ProcessContext context, final Process } }); - session.getProvenanceReporter().receive(messageFlowfile, getTransitUri(topicPrefix, topicFilter)); + session.getProvenanceReporter().receive(messageFlowfile, getTransitUri(getSubscribedTopicsForProvenance())); session.transfer(messageFlowfile, REL_MESSAGE); session.commitAsync(); } @@ -607,7 +667,7 @@ private void transferQueueRecord(final ProcessContext context, final ProcessSess } session.putAllAttributes(flowFile, attributes); - session.getProvenanceReporter().receive(flowFile, getTransitUri(topicPrefix, topicFilter)); + session.getProvenanceReporter().receive(flowFile, getTransitUri(getSubscribedTopicsForProvenance())); session.transfer(flowFile, REL_MESSAGE); final int count = recordCount.get(); @@ -645,6 +705,18 @@ private String getTransitUri(String... appends) { return stringBuilder.toString(); } + /** + * Returns the subscribed topic filters, including the shared subscription prefix if any, as a single comma + * separated value. A FlowFile produced by the demarcator or record based code paths may aggregate messages of + * several topics, so the individual topic of a message cannot be used for those. For a single Topic Filter this + * returns exactly the previously reported value. + */ + private String getSubscribedTopicsForProvenance() { + return topicSubscriptions.stream() + .map(MqttTopicSubscription::topicFilter) + .collect(Collectors.joining(",")); + } + private void handleReceivedMessage(ReceivedMqttMessage message) { if (logger.isDebugEnabled()) { byte[] payload = message.getPayload(); diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/HiveMqV5ClientAdapter.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/HiveMqV5ClientAdapter.java index 9aa927f282de..b703ada6ae2e 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/HiveMqV5ClientAdapter.java +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/HiveMqV5ClientAdapter.java @@ -22,12 +22,16 @@ import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientBuilder; import com.hivemq.client.mqtt.mqtt5.message.connect.Mqtt5Connect; import com.hivemq.client.mqtt.mqtt5.message.connect.Mqtt5ConnectBuilder; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.Mqtt5Subscribe; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.Mqtt5Subscription; import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAck; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAckReasonCode; import org.apache.nifi.logging.ComponentLog; import org.apache.nifi.processors.mqtt.common.MqttClient; import org.apache.nifi.processors.mqtt.common.MqttClientProperties; import org.apache.nifi.processors.mqtt.common.MqttException; import org.apache.nifi.processors.mqtt.common.MqttProtocolScheme; +import org.apache.nifi.processors.mqtt.common.MqttTopicSubscription; import org.apache.nifi.processors.mqtt.common.ReceivedMqttMessage; import org.apache.nifi.processors.mqtt.common.ReceivedMqttMessageHandler; import org.apache.nifi.processors.mqtt.common.StandardMqttMessage; @@ -36,10 +40,13 @@ import java.net.URI; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; import java.util.Objects; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; import javax.net.ssl.KeyManagerFactory; import javax.net.ssl.TrustManagerFactory; import javax.net.ssl.X509ExtendedKeyManager; @@ -61,6 +68,13 @@ public HiveMqV5ClientAdapter(URI brokerUri, MqttClientProperties clientPropertie this.logger = logger; } + // Package-private constructor for injecting a test double for the underlying HiveMQ client. + HiveMqV5ClientAdapter(Mqtt5BlockingClient mqtt5BlockingClient, MqttClientProperties clientProperties, ComponentLog logger) { + this.mqtt5BlockingClient = mqtt5BlockingClient; + this.clientProperties = clientProperties; + this.logger = logger; + } + @Override public boolean isConnected() { return mqtt5BlockingClient.getState().isConnected(); @@ -127,30 +141,51 @@ public void publish(String topic, StandardMqttMessage message) { } @Override - public void subscribe(String topicFilter, int qos, ReceivedMqttMessageHandler handler) { - logger.debug("Subscribing to {} with QoS: {}", topicFilter, qos); - - CompletableFuture futureAck = mqtt5BlockingClient.toAsync().subscribeWith() - .topicFilter(topicFilter) - .qos(Objects.requireNonNull(MqttQos.fromCode(qos))) - .callback(mqtt5Publish -> { - final ReceivedMqttMessage receivedMessage = new ReceivedMqttMessage( - mqtt5Publish.getPayloadAsBytes(), - mqtt5Publish.getQos().getCode(), - mqtt5Publish.isRetain(), - mqtt5Publish.getTopic().toString()); - handler.handleReceivedMessage(receivedMessage); - }) - .send(); - - // Setting "listener" callback is only possible with async client, though sending subscribe message - // should happen in a blocking way to make sure the processor is blocked until ack is not arrived. + public void subscribe(List subscriptions, ReceivedMqttMessageHandler handler) { + logger.debug("Subscribing to {}", subscriptions); + + final List mqtt5Subscriptions = subscriptions.stream() + .map(subscription -> Mqtt5Subscription.builder() + .topicFilter(subscription.topicFilter()) + .qos(Objects.requireNonNull(MqttQos.fromCode(subscription.qos()))) + .build()) + .collect(Collectors.toList()); + + final Mqtt5Subscribe mqtt5Subscribe = Mqtt5Subscribe.builder() + .addSubscriptions(mqtt5Subscriptions) + .build(); + + // Setting the "listener" callback is only possible with the async client, though sending the subscribe + // message should happen in a blocking way to make sure the processor is blocked until the ack arrives. + final CompletableFuture futureAck = mqtt5BlockingClient.toAsync().subscribe(mqtt5Subscribe, mqtt5Publish -> { + final ReceivedMqttMessage receivedMessage = new ReceivedMqttMessage( + mqtt5Publish.getPayloadAsBytes(), + mqtt5Publish.getQos().getCode(), + mqtt5Publish.isRetain(), + mqtt5Publish.getTopic().toString()); + handler.handleReceivedMessage(receivedMessage); + }); + + final Mqtt5SubAck ack; try { - final Mqtt5SubAck ack = futureAck.get(clientProperties.getConnectionTimeout(), TimeUnit.SECONDS); - logger.debug("Received mqtt5 subscribe ack: {}", ack); + ack = futureAck.get(clientProperties.getConnectionTimeout(), TimeUnit.SECONDS); } catch (Exception e) { throw new MqttException("An error has occurred during sending subscribe message to broker", e); } + logger.debug("Received mqtt5 subscribe ack: {}", ack); + + // A SUBACK carries one reason code per requested Topic Filter, in the order they were sent, so a subscription + // can be rejected individually, for example due to an ACL denial, while the others are granted. + final List reasonCodes = ack.getReasonCodes(); + final List failedTopicFilters = new ArrayList<>(); + for (int i = 0; i < reasonCodes.size() && i < subscriptions.size(); i++) { + if (reasonCodes.get(i).isError()) { + failedTopicFilters.add(subscriptions.get(i).topicFilter() + " (" + reasonCodes.get(i) + ")"); + } + } + if (!failedTopicFilters.isEmpty()) { + throw new MqttException("Broker rejected subscription for the following topic filter(s): " + failedTopicFilters); + } } private static Mqtt5BlockingClient createClient(URI brokerUri, MqttClientProperties clientProperties, ComponentLog logger) throws TlsException { diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/PahoMqttClientAdapter.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/PahoMqttClientAdapter.java index ddef128eb94d..259c6cff4ea1 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/PahoMqttClientAdapter.java +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/adapters/PahoMqttClientAdapter.java @@ -20,24 +20,32 @@ import org.apache.nifi.processors.mqtt.common.MqttClient; import org.apache.nifi.processors.mqtt.common.MqttClientProperties; import org.apache.nifi.processors.mqtt.common.MqttException; +import org.apache.nifi.processors.mqtt.common.MqttTopicSubscription; import org.apache.nifi.processors.mqtt.common.ReceivedMqttMessage; import org.apache.nifi.processors.mqtt.common.ReceivedMqttMessageHandler; import org.apache.nifi.processors.mqtt.common.StandardMqttMessage; import org.apache.nifi.ssl.SSLContextProvider; import org.eclipse.paho.client.mqttv3.IMqttClient; import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken; +import org.eclipse.paho.client.mqttv3.IMqttToken; import org.eclipse.paho.client.mqttv3.MqttCallback; import org.eclipse.paho.client.mqttv3.MqttConnectOptions; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import java.net.URI; +import java.util.ArrayList; import java.util.Arrays; +import java.util.List; public class PahoMqttClientAdapter implements MqttClient { public static final int DISCONNECT_TIMEOUT = 5000; + // MQTT SUBACK reason code indicating the broker rejected the requested subscription, for example due to an ACL + // denial. Paho does not surface this as an error on its own, so the granted QoS array has to be checked here. + private static final int SUBACK_FAILURE_CODE = 0x80; + private final IMqttClient client; private final MqttClientProperties clientProperties; private final ComponentLog logger; @@ -49,6 +57,14 @@ public PahoMqttClientAdapter(URI brokerUri, MqttClientProperties clientPropertie client.setCallback(new DefaultMqttCallback()); } + // Package-private constructor for injecting a test double for the underlying Paho client. + PahoMqttClientAdapter(IMqttClient client, MqttClientProperties clientProperties, ComponentLog logger) { + this.client = client; + this.clientProperties = clientProperties; + this.logger = logger; + client.setCallback(new DefaultMqttCallback()); + } + @Override public boolean isConnected() { return client.isConnected(); @@ -123,15 +139,28 @@ public void publish(String topic, StandardMqttMessage message) { } @Override - public void subscribe(String topicFilter, int qos, ReceivedMqttMessageHandler handler) { - logger.debug("Subscribing to {} with QoS: {}", topicFilter, qos); + public void subscribe(List subscriptions, ReceivedMqttMessageHandler handler) { + final String[] topicFilters = subscriptions.stream().map(MqttTopicSubscription::topicFilter).toArray(String[]::new); + final int[] qosLevels = subscriptions.stream().mapToInt(MqttTopicSubscription::qos).toArray(); + + logger.debug("Subscribing to {} with QoS: {}", Arrays.toString(topicFilters), Arrays.toString(qosLevels)); client.setCallback(new ConsumerMqttCallback(handler)); try { - client.subscribe(topicFilter, qos); + final IMqttToken token = client.subscribeWithResponse(topicFilters, qosLevels); + final int[] grantedQos = token.getGrantedQos(); + final List failedTopicFilters = new ArrayList<>(); + for (int i = 0; i < grantedQos.length; i++) { + if (grantedQos[i] == SUBACK_FAILURE_CODE) { + failedTopicFilters.add(topicFilters[i]); + } + } + if (!failedTopicFilters.isEmpty()) { + throw new MqttException("Broker rejected subscription for the following topic filter(s): " + failedTopicFilters); + } } catch (org.eclipse.paho.client.mqttv3.MqttException e) { - throw new MqttException("An error has occurred during subscribing to " + topicFilter + " with QoS: " + qos, e); + throw new MqttException("An error has occurred during subscribing to " + Arrays.toString(topicFilters) + " with QoS: " + Arrays.toString(qosLevels), e); } } diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttClient.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttClient.java index 2b5c949531d0..3dfeb99c4561 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttClient.java +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttClient.java @@ -16,6 +16,8 @@ */ package org.apache.nifi.processors.mqtt.common; +import java.util.List; + public interface MqttClient { /** @@ -50,14 +52,14 @@ public interface MqttClient { void publish(String topic, StandardMqttMessage message); /** - * Subscribe to a topic. + * Subscribe to one or more topic filters with a single SUBSCRIBE request. * - * @param topicFilter the topic to subscribe to, which can include wildcards. - * @param qos the maximum quality of service at which to subscribe. Messages - * published at a lower quality of service will be received at the published - * QoS. Messages published at a higher quality of service will be received using - * the QoS specified on the subscribe. + * @param subscriptions the list of (Topic Filter, QoS) pairs to subscribe to. A topic filter can include + * wildcards. The QoS is the maximum quality of service at which to subscribe. Messages + * published at a lower quality of service will be received at the published QoS. Messages + * published at a higher quality of service will be received using the QoS specified on + * the subscribe. * @param handler that further processes the message received by the client */ - void subscribe(String topicFilter, int qos, ReceivedMqttMessageHandler handler); + void subscribe(List subscriptions, ReceivedMqttMessageHandler handler); } diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttTopicSubscription.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttTopicSubscription.java new file mode 100644 index 000000000000..582d9e631f74 --- /dev/null +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/java/org/apache/nifi/processors/mqtt/common/MqttTopicSubscription.java @@ -0,0 +1,29 @@ +/* + * 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.nifi.processors.mqtt.common; + +/** + * Represents a single (Topic Filter, QoS) pair of a SUBSCRIBE request. A SUBSCRIBE packet can carry a list of these, + * allowing a client to subscribe to multiple Topic Filters with a single request. + * + * @param topicFilter the topic filter to subscribe to, which can include wildcards. + * @param qos the maximum quality of service at which to subscribe. Messages published at a lower quality of + * service will be received at the published QoS. Messages published at a higher quality of service + * will be received using the QoS specified on the subscribe. + */ +public record MqttTopicSubscription(String topicFilter, int qos) { +} diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/resources/docs/org.apache.nifi.processors.mqtt.ConsumeMQTT/additionalDetails.md b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/resources/docs/org.apache.nifi.processors.mqtt.ConsumeMQTT/additionalDetails.md index bec36213f34b..3b92ed4e213f 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/resources/docs/org.apache.nifi.processors.mqtt.ConsumeMQTT/additionalDetails.md +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/main/resources/docs/org.apache.nifi.processors.mqtt.ConsumeMQTT/additionalDetails.md @@ -22,4 +22,26 @@ messages in the internal queue will be written to FlowFiles. In case the interna try for up to 1 second to add the message into the internal queue. If the internal queue is still full after this time, an exception saying that 'The subscriber queue is full' would be thrown, the message would be dropped and the client would be disconnected. In case the QoS property is set to 0, the message would be lost. In case the QoS property is set -to 1 or 2, the message will be received after the client reconnects. \ No newline at end of file +to 1 or 2, the message will be received after the client reconnects. + +## Multiple Topic Filters + +The 'Topic Filter' property accepts a comma-separated list of topic filters. When more than one filter is provided, all +of them are subscribed to with a single SUBSCRIBE request over the same broker connection, instead of requiring one +instance of this processor (and one broker connection) per topic. This is especially useful when the topics to consume +are flat or externally dictated, or when broker ACLs only authorize specific, explicit topics, making wildcard filters +unusable. + +Each topic filter in the list is subscribed to at the same Quality of Service, configured by the 'Quality of Service' +property. If 'Group ID' is set, the shared subscription prefix (`$share//`) is applied to every topic filter +individually, so each remains its own subscription eligible for load-balancing across the group. + +A value that contains no comma is a single topic filter and is used exactly as configured, so existing configurations +are unaffected. Within a comma-separated list, whitespace around each topic filter is trimmed, empty entries are +ignored, and duplicate topic filters are rejected during validation. Because MQTT topic filters may legally contain a +comma, a topic filter that itself contains a comma will be misinterpreted as multiple filters; this is expected to be a +rare edge case. + +If the broker rejects a subscription for one or more of the requested topic filters (for example due to an ACL denial), +the processor fails to initialize its client and yields, logging the offending topic filter(s), rather than silently +continuing without them. diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/TestConsumeMQTT.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/TestConsumeMQTT.java index ce3eed059e0a..ccba745bef8e 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/TestConsumeMQTT.java +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/TestConsumeMQTT.java @@ -21,7 +21,9 @@ import org.apache.nifi.processor.ProcessSession; import org.apache.nifi.processors.mqtt.common.AbstractMQTTProcessor; import org.apache.nifi.processors.mqtt.common.MqttClient; +import org.apache.nifi.processors.mqtt.common.MqttException; import org.apache.nifi.processors.mqtt.common.MqttTestClient; +import org.apache.nifi.processors.mqtt.common.MqttTopicSubscription; import org.apache.nifi.processors.mqtt.common.ReceivedMqttMessage; import org.apache.nifi.processors.mqtt.common.StandardMqttMessage; import org.apache.nifi.provenance.ProvenanceEventRecord; @@ -44,6 +46,7 @@ import java.util.List; import java.util.Map; import java.util.concurrent.BlockingQueue; +import java.util.function.Consumer; import static org.apache.nifi.processors.mqtt.ConsumeMQTT.BROKER_ATTRIBUTE_KEY; import static org.apache.nifi.processors.mqtt.ConsumeMQTT.IS_DUPLICATE_ATTRIBUTE_KEY; @@ -109,6 +112,166 @@ public void testClientIDConfiguration() { testRunner.assertValid(); } + @Test + public void testSingleTopicFilterSubscription() throws Exception { + final List subscriptions = subscribeAndGetSubscriptions(runner -> { }); + + assertEquals(List.of(new MqttTopicSubscription(TOPIC_NAME, AT_MOST_ONCE)), subscriptions); + } + + @Test + public void testSingleTopicFilterSubscriptionIsUsedVerbatim() throws Exception { + // Leading and trailing whitespace is significant in an MQTT topic filter, so a value without a comma must not + // be altered, otherwise upgrading would silently change the subscription of an existing configuration. + final List subscriptions = subscribeAndGetSubscriptions( + runner -> runner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, " testTopic/one ")); + + assertEquals(List.of(new MqttTopicSubscription(" testTopic/one ", AT_MOST_ONCE)), subscriptions); + } + + @Test + public void testMultipleTopicFilterSubscription() throws Exception { + final List subscriptions = subscribeAndGetSubscriptions( + runner -> runner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one,testTopic/two,testTopic/three")); + + assertEquals(List.of( + new MqttTopicSubscription("testTopic/one", AT_MOST_ONCE), + new MqttTopicSubscription("testTopic/two", AT_MOST_ONCE), + new MqttTopicSubscription("testTopic/three", AT_MOST_ONCE)), subscriptions); + } + + @Test + public void testMultipleTopicFilterSubscriptionTrimsWhitespaceAndIgnoresEmptyEntries() throws Exception { + final List subscriptions = subscribeAndGetSubscriptions( + runner -> runner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, " testTopic/one ,, testTopic/two ,")); + + assertEquals(List.of( + new MqttTopicSubscription("testTopic/one", AT_MOST_ONCE), + new MqttTopicSubscription("testTopic/two", AT_MOST_ONCE)), subscriptions); + } + + @Test + public void testMultipleTopicFilterSubscriptionUsesConfiguredQos() throws Exception { + final List subscriptions = subscribeAndGetSubscriptions(runner -> { + runner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one,testTopic/two"); + runner.setProperty(ConsumeMQTT.PROP_QOS, String.valueOf(EXACTLY_ONCE)); + }); + + assertEquals(List.of( + new MqttTopicSubscription("testTopic/one", EXACTLY_ONCE), + new MqttTopicSubscription("testTopic/two", EXACTLY_ONCE)), subscriptions); + } + + @Test + public void testMultipleTopicFilterSubscriptionAppliesSharedSubscriptionPrefixPerFilter() throws Exception { + final List subscriptions = subscribeAndGetSubscriptions(runner -> { + runner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one,testTopic/two"); + runner.setProperty(ConsumeMQTT.PROP_GROUPID, "testGroup"); + runner.setProperty(ConsumeMQTT.PROP_CLIENTID, "${hostname()}"); + }); + + assertEquals(List.of( + new MqttTopicSubscription("$share/testGroup/testTopic/one", AT_MOST_ONCE), + new MqttTopicSubscription("$share/testGroup/testTopic/two", AT_MOST_ONCE)), subscriptions); + } + + @Test + public void testTopicFilterWithExpressionLanguageIsValid() { + mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); + testRunner = initializeTestRunner(mqttTestClient); + + testRunner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "${literal('testTopic/one,testTopic/two')}"); + testRunner.assertValid(); + } + + @Test + public void testDuplicateTopicFiltersNotValid() { + mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); + testRunner = initializeTestRunner(mqttTestClient); + + testRunner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one,testTopic/two"); + testRunner.assertValid(); + + testRunner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one, testTopic/two , testTopic/one"); + testRunner.assertNotValid(); + } + + @Test + public void testBlankTopicFiltersNotValid() { + mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); + testRunner = initializeTestRunner(mqttTestClient); + + testRunner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, " , , "); + testRunner.assertNotValid(); + } + + @Test + public void testClientIsDisconnectedAndClosedWhenSubscribeFails() throws Exception { + mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); + mqttTestClient.subscribeException = new MqttException("Broker rejected subscription for the following topic filter(s): [testTopic/two]"); + testRunner = initializeTestRunner(mqttTestClient); + testRunner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one,testTopic/two"); + + final ConsumeMQTT consumeMQTT = (ConsumeMQTT) testRunner.getProcessor(); + consumeMQTT.onScheduled(testRunner.getProcessContext()); + reconnect(consumeMQTT, testRunner.getProcessContext()); + + // A SUBSCRIBE listing several topic filters can be granted partially, leaving the client connected and + // subscribed. Dereferencing it without disconnecting would leak the broker connection and let the orphaned + // client keep feeding the internal queue while retries create further clients. + assertFalse(mqttTestClient.connected.get()); + assertTrue(mqttTestClient.closed.get()); + } + + @Test + public void testMultipleTopicFilterProvenanceTransitUriForAggregatedFlowFile() throws Exception { + mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); + testRunner = initializeTestRunner(mqttTestClient); + + testRunner.setProperty(ConsumeMQTT.PROP_TOPIC_FILTER, "testTopic/one,testTopic/two"); + testRunner.setProperty(ConsumeMQTT.PROP_GROUPID, "testGroup"); + testRunner.setProperty(ConsumeMQTT.PROP_CLIENTID, "${hostname()}"); + testRunner.setProperty(ConsumeMQTT.RECORD_READER, createJsonRecordSetReaderService(testRunner)); + testRunner.setProperty(ConsumeMQTT.RECORD_WRITER, createJsonRecordSetWriterService(testRunner)); + + testRunner.assertValid(); + + final ConsumeMQTT consumeMQTT = (ConsumeMQTT) testRunner.getProcessor(); + consumeMQTT.onScheduled(testRunner.getProcessContext()); + reconnect(consumeMQTT, testRunner.getProcessContext()); + + Thread.sleep(PUBLISH_WAIT_MS); + assertTrue(isConnected(consumeMQTT)); + + mqttTestClient.publish("testTopic/one", new StandardMqttMessage(JSON_PAYLOAD.getBytes(StandardCharsets.UTF_8), AT_MOST_ONCE, false)); + mqttTestClient.publish("testTopic/two", new StandardMqttMessage(JSON_PAYLOAD.getBytes(StandardCharsets.UTF_8), AT_MOST_ONCE, false)); + + Thread.sleep(PUBLISH_WAIT_MS); + + testRunner.run(1, false, false); + + // A single FlowFile aggregates messages of both topics, so the transit URI has to list every subscribed topic + // filter, each carrying the shared subscription prefix. + final List provenanceEvents = testRunner.getProvenanceEvents(); + assertEquals(1, provenanceEvents.size()); + assertEquals(BROKER_URI + "/$share/testGroup/testTopic/one,$share/testGroup/testTopic/two", + provenanceEvents.getFirst().getTransitUri()); + } + + private List subscribeAndGetSubscriptions(final Consumer configurer) throws Exception { + mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); + testRunner = initializeTestRunner(mqttTestClient); + + configurer.accept(testRunner); + testRunner.assertValid(); + + final ConsumeMQTT consumeMQTT = (ConsumeMQTT) testRunner.getProcessor(); + consumeMQTT.onScheduled(testRunner.getProcessContext()); + reconnect(consumeMQTT, testRunner.getProcessContext()); + + return mqttTestClient.subscriptions; + } + @Test public void testLastWillConfig() { mqttTestClient = new MqttTestClient(MqttTestClient.ConnectType.Subscriber); diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/adapters/TestHiveMqV5ClientAdapter.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/adapters/TestHiveMqV5ClientAdapter.java new file mode 100644 index 000000000000..72857e237c4a --- /dev/null +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/adapters/TestHiveMqV5ClientAdapter.java @@ -0,0 +1,96 @@ +/* + * 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.nifi.processors.mqtt.adapters; + +import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; +import com.hivemq.client.mqtt.mqtt5.Mqtt5BlockingClient; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.Mqtt5Subscribe; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAck; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAckReasonCode; +import org.apache.nifi.logging.ComponentLog; +import org.apache.nifi.processors.mqtt.common.MqttClientProperties; +import org.apache.nifi.processors.mqtt.common.MqttException; +import org.apache.nifi.processors.mqtt.common.MqttTopicSubscription; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.stream.Collectors; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +public class TestHiveMqV5ClientAdapter { + + @Test + public void testSubscribeMultipleTopicFiltersSendsSingleSubscribeRequest() { + final Mqtt5BlockingClient blockingClient = mock(Mqtt5BlockingClient.class); + final Mqtt5AsyncClient asyncClient = mock(Mqtt5AsyncClient.class); + when(blockingClient.toAsync()).thenReturn(asyncClient); + + final Mqtt5SubAck subAck = mock(Mqtt5SubAck.class); + when(subAck.getReasonCodes()).thenReturn(List.of(Mqtt5SubAckReasonCode.GRANTED_QOS_1, Mqtt5SubAckReasonCode.GRANTED_QOS_1)); + when(asyncClient.subscribe(any(Mqtt5Subscribe.class), any())).thenReturn(CompletableFuture.completedFuture(subAck)); + + final MqttClientProperties clientProperties = new MqttClientProperties(); + clientProperties.setConnectionTimeout(5); + final HiveMqV5ClientAdapter adapter = new HiveMqV5ClientAdapter(blockingClient, clientProperties, mock(ComponentLog.class)); + + final List subscriptions = List.of( + new MqttTopicSubscription("topic/a", 1), + new MqttTopicSubscription("topic/b", 1)); + + adapter.subscribe(subscriptions, message -> { }); + + // Both topic filters have to be carried by a single SUBSCRIBE request. + final ArgumentCaptor captor = ArgumentCaptor.forClass(Mqtt5Subscribe.class); + verify(asyncClient, times(1)).subscribe(captor.capture(), any()); + assertEquals(List.of("topic/a", "topic/b"), captor.getValue().getSubscriptions().stream() + .map(subscription -> subscription.getTopicFilter().toString()) + .collect(Collectors.toList())); + } + + @Test + public void testSubscribeThrowsWhenBrokerRejectsATopicFilter() { + final Mqtt5BlockingClient blockingClient = mock(Mqtt5BlockingClient.class); + final Mqtt5AsyncClient asyncClient = mock(Mqtt5AsyncClient.class); + when(blockingClient.toAsync()).thenReturn(asyncClient); + + final Mqtt5SubAck subAck = mock(Mqtt5SubAck.class); + // Broker grants the first filter but rejects the second, e.g. due to an ACL denial. + when(subAck.getReasonCodes()).thenReturn(List.of(Mqtt5SubAckReasonCode.GRANTED_QOS_1, Mqtt5SubAckReasonCode.NOT_AUTHORIZED)); + when(asyncClient.subscribe(any(Mqtt5Subscribe.class), any())).thenReturn(CompletableFuture.completedFuture(subAck)); + + final MqttClientProperties clientProperties = new MqttClientProperties(); + clientProperties.setConnectionTimeout(5); + final HiveMqV5ClientAdapter adapter = new HiveMqV5ClientAdapter(blockingClient, clientProperties, mock(ComponentLog.class)); + + final List subscriptions = List.of( + new MqttTopicSubscription("topic/a", 1), + new MqttTopicSubscription("topic/denied", 1)); + + final MqttException e = assertThrows(MqttException.class, () -> adapter.subscribe(subscriptions, message -> { })); + assertTrue(e.getMessage().contains("topic/denied")); + } +} diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/adapters/TestPahoMqttClientAdapter.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/adapters/TestPahoMqttClientAdapter.java new file mode 100644 index 000000000000..25fdd0df08ed --- /dev/null +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/adapters/TestPahoMqttClientAdapter.java @@ -0,0 +1,76 @@ +/* + * 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.nifi.processors.mqtt.adapters; + +import org.apache.nifi.logging.ComponentLog; +import org.apache.nifi.processors.mqtt.common.MqttClientProperties; +import org.apache.nifi.processors.mqtt.common.MqttException; +import org.apache.nifi.processors.mqtt.common.MqttTopicSubscription; +import org.eclipse.paho.client.mqttv3.IMqttClient; +import org.eclipse.paho.client.mqttv3.IMqttToken; +import org.junit.jupiter.api.Test; + +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +public class TestPahoMqttClientAdapter { + + private static final int SUBACK_FAILURE_CODE = 0x80; + + @Test + public void testSubscribeMultipleTopicFiltersSendsSingleSubscribeRequest() throws Exception { + final IMqttClient client = mock(IMqttClient.class); + final IMqttToken token = mock(IMqttToken.class); + when(client.subscribeWithResponse(any(String[].class), any(int[].class))).thenReturn(token); + when(token.getGrantedQos()).thenReturn(new int[]{1, 1}); + + final PahoMqttClientAdapter adapter = new PahoMqttClientAdapter(client, new MqttClientProperties(), mock(ComponentLog.class)); + + final List subscriptions = List.of( + new MqttTopicSubscription("topic/a", 1), + new MqttTopicSubscription("topic/b", 1)); + + adapter.subscribe(subscriptions, message -> { }); + + verify(client, times(1)).subscribeWithResponse(new String[]{"topic/a", "topic/b"}, new int[]{1, 1}); + } + + @Test + public void testSubscribeThrowsWhenBrokerRejectsATopicFilter() throws Exception { + final IMqttClient client = mock(IMqttClient.class); + final IMqttToken token = mock(IMqttToken.class); + when(client.subscribeWithResponse(any(String[].class), any(int[].class))).thenReturn(token); + // Broker grants the first filter but rejects the second, e.g. due to an ACL denial. + when(token.getGrantedQos()).thenReturn(new int[]{1, SUBACK_FAILURE_CODE}); + + final PahoMqttClientAdapter adapter = new PahoMqttClientAdapter(client, new MqttClientProperties(), mock(ComponentLog.class)); + + final List subscriptions = List.of( + new MqttTopicSubscription("topic/a", 1), + new MqttTopicSubscription("topic/denied", 1)); + + final MqttException e = assertThrows(MqttException.class, () -> adapter.subscribe(subscriptions, message -> { })); + assertTrue(e.getMessage().contains("topic/denied")); + } +} diff --git a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/common/MqttTestClient.java b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/common/MqttTestClient.java index 697a80da89de..23bbe384949a 100644 --- a/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/common/MqttTestClient.java +++ b/nifi-extension-bundles/nifi-mqtt-bundle/nifi-mqtt-processors/src/test/java/org/apache/nifi/processors/mqtt/common/MqttTestClient.java @@ -20,6 +20,7 @@ import org.apache.commons.lang3.tuple.Pair; import java.util.LinkedList; +import java.util.List; import java.util.Queue; import java.util.concurrent.atomic.AtomicBoolean; @@ -29,12 +30,16 @@ public class MqttTestClient implements MqttClient { public AtomicBoolean connected = new AtomicBoolean(false); + public AtomicBoolean closed = new AtomicBoolean(false); + + // Allows simulating a broker rejecting the SUBSCRIBE after the client has already connected. + public RuntimeException subscribeException; + public ConnectType type; public enum ConnectType { Publisher, Subscriber } - public String subscribedTopic; - public int subscribedQos; + public List subscriptions; public ReceivedMqttMessageHandler receivedMqttMessageHandler; public MqttTestClient(ConnectType type) { this.type = type; @@ -57,7 +62,7 @@ public void disconnect() { @Override public void close() { - + closed.set(true); } @Override @@ -73,10 +78,12 @@ public void publish(String topic, StandardMqttMessage message) { } @Override - public void subscribe(String topicFilter, int qos, ReceivedMqttMessageHandler handler) { - subscribedTopic = topicFilter; - subscribedQos = qos; + public void subscribe(List subscriptions, ReceivedMqttMessageHandler handler) { + this.subscriptions = subscriptions; receivedMqttMessageHandler = handler; + if (subscribeException != null) { + throw subscribeException; + } } public Pair getLastPublished() {