diff --git a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ConnectivityModelFactory.java b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ConnectivityModelFactory.java index cd642052cf2..9837bf4dab9 100755 --- a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ConnectivityModelFactory.java +++ b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ConnectivityModelFactory.java @@ -842,8 +842,9 @@ public static Target targetFromJson(final JsonObject jsonObject) { } /** - * Creates a new {@code FilteredTopic} from the passed {@code topicString} which consists of a {@code Topic} and an - * optional filter string supplied with {@code ?filter=...}. + * Creates a new {@code FilteredTopic} from the passed {@code topicString} which consists of a {@code Topic} and + * optional filter strings supplied with {@code ?filter=...}. The {@code filter} query parameter may be repeated + * ({@code ?filter=...&filter=...}); all given filters must match for a signal to be processed (AND semantics). * * @param topicString the {@code FilteredTopic} String representation * @return the created FilteredTopic diff --git a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopic.java b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopic.java index 5e052a33e03..19ee76208ab 100644 --- a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopic.java +++ b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopic.java @@ -18,8 +18,10 @@ import org.eclipse.ditto.json.JsonFieldSelector; /** - * A FilteredTopic wraps a {@link Topic} and an optional {@code filter} String which additionally restricts which - * kind of Signals should be processed/filtered based on an {@code RQL} query. + * A FilteredTopic wraps a {@link Topic} and optional {@code filter} Strings which additionally restrict which + * kind of Signals should be processed/filtered. Each filter is either an {@code RQL} query or a placeholder + * pipeline expression starting with {@code fn:}; all filters of one topic must match for a signal to be + * processed (AND semantics). */ public interface FilteredTopic extends CharSequence { @@ -34,10 +36,22 @@ public interface FilteredTopic extends CharSequence { List getNamespaces(); /** - * @return the optional filter string as RQL query + * @return the first filter string of this FilteredTopic, or an empty Optional if no filter is set. + * @deprecated as of 3.10.0 a FilteredTopic may carry multiple filters; use {@link #getFilters()} instead. */ + @Deprecated Optional getFilter(); + /** + * Returns the filter strings of this FilteredTopic in insertion order. All filters of one topic must match + * for a signal to be processed (AND semantics). At most one entry may be an RQL expression; any number of + * entries may be placeholder pipeline expressions starting with {@code fn:}. + * + * @return the filter strings, or an empty list if no filter is set. + * @since 3.10.0 + */ + List getFilters(); + /** * Returns the selector for the extra fields and their values to enrich outgoing signals with. * diff --git a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopicBuilder.java b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopicBuilder.java index 105555a27af..7879f46f045 100644 --- a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopicBuilder.java +++ b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/FilteredTopicBuilder.java @@ -35,13 +35,28 @@ public interface FilteredTopicBuilder { FilteredTopicBuilder withNamespaces(@Nullable Collection namespaces); /** - * Sets the given filter to this builder. + * Sets the given filter to this builder, replacing all previously set filters. * - * @param filter the optional RQL filter of the topic to be built. + * @param filter the optional filter of the topic to be built. * @return this builder instance to allow method chaining. + * @deprecated as of 3.10.0 a FilteredTopic may carry multiple filters; use {@link #withFilters(Collection)} + * instead. */ + @Deprecated FilteredTopicBuilder withFilter(@Nullable CharSequence filter); + /** + * Sets the given filters to this builder, replacing all previously set filters. The insertion order is + * preserved and determines the serialization order of the {@code filter} query parameters; two topics with + * the same filters in different order are not equal. + * + * @param filters the filters of the topic to be built - each entry is either an RQL expression or a + * placeholder pipeline expression starting with {@code fn:}. + * @return this builder instance to allow method chaining. + * @since 3.10.0 + */ + FilteredTopicBuilder withFilters(@Nullable Collection filters); + /** * Sets the selector for the extra fields and their values to enrich outgoing signals of the topic to be built with. * diff --git a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopic.java b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopic.java index b2b67fae992..b340e4f1faf 100644 --- a/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopic.java +++ b/connectivity/model/src/main/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopic.java @@ -21,6 +21,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.Collections; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; @@ -51,7 +52,7 @@ final class ImmutableFilteredTopic implements FilteredTopic { private final Topic topic; private final List namespaces; - @Nullable private final String filterString; + private final List filters; @Nullable private final ThingFieldSelector extraFields; private ImmutableFilteredTopic(final ImmutableFilteredTopicBuilder builder) { @@ -60,7 +61,10 @@ private ImmutableFilteredTopic(final ImmutableFilteredTopicBuilder builder) { namespaces = null != namespacesFromBuilder ? Collections.unmodifiableList(new ArrayList<>(namespacesFromBuilder)) : Collections.emptyList(); - filterString = Objects.toString(builder.filter, null); + final Collection filtersFromBuilder = builder.filters; + filters = null != filtersFromBuilder + ? Collections.unmodifiableList(new ArrayList<>(filtersFromBuilder)) + : Collections.emptyList(); extraFields = builder.extraFields; } @@ -102,7 +106,12 @@ public List getNamespaces() { @Override public Optional getFilter() { - return Optional.ofNullable(filterString); + return filters.isEmpty() ? Optional.empty() : Optional.of(filters.get(0)); + } + + @Override + public List getFilters() { + return filters; } @Override @@ -131,9 +140,13 @@ public String toString() { } private String getQueryParametersAsString() { - return join(QUERY_ARG_DELIMITER, getQueryParameterString(NAMESPACES_ARG, String.join(",", namespaces)), - getQueryParameterString(FILTER_ARG, filterString), - getQueryParameterString(EXTRA_FIELDS_ARG, extraFields)); + final List queryParameterStrings = new ArrayList<>(2 + filters.size()); + queryParameterStrings.add(getQueryParameterString(NAMESPACES_ARG, String.join(",", namespaces))); + for (final String filter : filters) { + queryParameterStrings.add(getQueryParameterString(FILTER_ARG, filter)); + } + queryParameterStrings.add(getQueryParameterString(EXTRA_FIELDS_ARG, extraFields)); + return join(QUERY_ARG_DELIMITER, queryParameterStrings.toArray(new String[0])); } private static String getQueryParameterString(final String parameterName, @Nullable final Object parameterValue) { @@ -169,13 +182,13 @@ public boolean equals(@Nullable final Object o) { final ImmutableFilteredTopic that = (ImmutableFilteredTopic) o; return topic == that.topic && namespaces.equals(that.namespaces) && - Objects.equals(filterString, that.filterString) && + filters.equals(that.filters) && Objects.equals(extraFields, that.extraFields); } @Override public int hashCode() { - return Objects.hash(topic, namespaces, filterString, extraFields); + return Objects.hash(topic, namespaces, filters, extraFields); } /** @@ -186,13 +199,13 @@ static final class ImmutableFilteredTopicBuilder implements FilteredTopicBuilder private final Topic topic; @Nullable private Collection namespaces; - @Nullable private CharSequence filter; + @Nullable private List filters; @Nullable private ThingFieldSelector extraFields; private ImmutableFilteredTopicBuilder(final Topic topic) { this.topic = checkNotNull(topic, "topic"); namespaces = null; - filter = null; + filters = null; extraFields = null; } @@ -206,8 +219,15 @@ public ImmutableFilteredTopicBuilder withNamespaces(@Nullable final Collection filters) { if (supportsFilters()) { - this.filter = filter; + this.filters = null != filters + ? filters.stream().map(CharSequence::toString).collect(Collectors.toList()) + : null; } return this; } @@ -239,13 +259,15 @@ private boolean supportsExtraFields() { } - @Immutable + @NotThreadSafe private static final class FilteredTopicStringParser { private final String filteredTopicString; + private final List filterValues; private FilteredTopicStringParser(final String filteredTopicString) { this.filteredTopicString = filteredTopicString; + filterValues = new ArrayList<>(1); } ImmutableFilteredTopic parse() { @@ -260,7 +282,7 @@ ImmutableFilteredTopic parse() { return getBuilder(parseTopic(topicName)) .withNamespaces(parseNamespaces(queryParameters.get(NAMESPACES_ARG))) - .withFilter(queryParameters.get(FILTER_ARG)) + .withFilters(filterValues.isEmpty() ? null : filterValues) .withExtraFields(parseExtraFields(queryParameters.get(EXTRA_FIELDS_ARG))) .build(); } @@ -271,14 +293,26 @@ private Topic parseTopic(final String topicName) { "Unknown topic: " + topicName).build()); } - private static Map parseQueryParameters(@Nullable final String queryParamsString) { + private Map parseQueryParameters(@Nullable final String queryParamsString) { if (null == queryParamsString || queryParamsString.isEmpty()) { return Collections.emptyMap(); } - return Arrays.stream(queryParamsString.split(QUERY_ARG_DELIMITER)) - .map(paramString -> paramString.split(QUERY_ARG_VALUE_DELIMITER, 2)) - .filter(queryParamPair -> 2 == queryParamPair.length) - .collect(Collectors.toMap(queryParamPair -> urlDecode(queryParamPair[0]), av -> urlDecode(av[1]))); + final Map queryParameters = new HashMap<>(4); + for (final String paramString : queryParamsString.split(QUERY_ARG_DELIMITER)) { + final String[] queryParamPair = paramString.split(QUERY_ARG_VALUE_DELIMITER, 2); + if (2 != queryParamPair.length) { + continue; + } + final String name = urlDecode(queryParamPair[0]); + final String value = urlDecode(queryParamPair[1]); + if (FILTER_ARG.equals(name)) { + // the filter parameter is the only repeatable one: values are collected in insertion order + filterValues.add(value); + } else if (null != queryParameters.putIfAbsent(name, value)) { + throw new IllegalStateException("Duplicate key " + name); + } + } + return queryParameters; } private static String urlDecode(final String value) { diff --git a/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopicTest.java b/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopicTest.java index 9b2ae10e62d..c947936cb23 100644 --- a/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopicTest.java +++ b/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableFilteredTopicTest.java @@ -15,6 +15,7 @@ import static org.assertj.core.api.Assertions.assertThat; import java.text.MessageFormat; +import java.util.Arrays; import java.util.Collections; import java.util.List; @@ -247,4 +248,119 @@ public void fromStringParsesAsExpectedWithOnlyExtraFields() { assertThat(actual).isEqualTo(filteredTopic); } + @Test + public void fromStringParsesAsExpectedWithNamespacesExtraFieldsAndRqlAndPipelineFilters() { + final String rqlFilter = "gt(attributes/counter,42)"; + final String pipelineFilter = "fn:filter(header:ditto-originator,'ne','some:subject')"; + final ImmutableFilteredTopic filteredTopic = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withNamespaces(Lists.list("ns1", "ns2")) + .withExtraFields(ThingFieldSelector.fromString("attributes")) + .withFilters(Arrays.asList(rqlFilter, pipelineFilter)) + .build(); + final String filteredTopicString = filteredTopic.toString(); + + final ImmutableFilteredTopic actual = ImmutableFilteredTopic.fromString(filteredTopicString); + + assertThat(filteredTopicString).contains("filter=" + rqlFilter + "&filter=" + pipelineFilter); + assertThat(actual.getFilters()).containsExactly(rqlFilter, pipelineFilter); + assertThat(actual.toString()).isEqualTo(filteredTopicString); + assertThat(actual).isEqualTo(filteredTopic); + } + + @Test + public void getFiltersReturnsAllInInsertionOrder() { + final ImmutableFilteredTopic underTest = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withFilters(Arrays.asList("fn:filter(header:b,'exists')", FILTER_EXAMPLE, "fn:filter(header:a,'exists')")) + .build(); + + assertThat(underTest.getFilters()) + .containsExactly("fn:filter(header:b,'exists')", FILTER_EXAMPLE, "fn:filter(header:a,'exists')"); + } + + @Test + public void getFilterReturnsFirstOfMultipleFilters() { + // contract of the deprecated single-filter accessor: the FIRST filter in insertion order + final ImmutableFilteredTopic underTest = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withFilters(Arrays.asList(FILTER_EXAMPLE, "fn:filter(header:a,'exists')")) + .build(); + + assertThat(underTest.getFilter()).contains(FILTER_EXAMPLE); + } + + @Test + public void withFilterReplacesPreviouslySetFilters() { + final ImmutableFilteredTopic underTest = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withFilters(Arrays.asList("fn:filter(header:a,'exists')", "fn:filter(header:b,'exists')")) + .withFilter(FILTER_EXAMPLE) + .build(); + + assertThat(underTest.getFilters()).containsExactly(FILTER_EXAMPLE); + } + + @Test + public void withFiltersNullResetsFilters() { + final ImmutableFilteredTopic underTest = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withFilters(Collections.singletonList(FILTER_EXAMPLE)) + .withFilters(null) + .build(); + + assertThat(underTest.getFilters()).isEmpty(); + assertThat(underTest.getFilter()).isEmpty(); + } + + @Test + public void toStringEmitsOneFilterParamPerEntryInOrder() { + final ImmutableFilteredTopic underTest = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withNamespaces(NAMESPACES) + .withFilters(Arrays.asList(FILTER_EXAMPLE, "fn:filter(header:a,'exists')")) + .withExtraFields(EXTRA_FIELDS) + .build(); + + assertThat(underTest.toString()).isEqualTo( + "_/_/things/twin/events?namespaces=" + String.join(",", NAMESPACES) + + "&filter=" + FILTER_EXAMPLE + + "&filter=fn:filter(header:a,'exists')" + + "&extraFields=" + EXTRA_FIELDS); + } + + @Test + public void fromStringCollectsRepeatedFilterParamsInOrder() { + final ImmutableFilteredTopic actual = ImmutableFilteredTopic.fromString( + "_/_/things/twin/events?filter=fn:filter(header:a,'exists')&filter=" + FILTER_EXAMPLE); + + assertThat(actual.getFilters()).containsExactly("fn:filter(header:a,'exists')", FILTER_EXAMPLE); + } + + @Test + public void fromStringToStringRoundTripsWithMultipleFilters() { + final ImmutableFilteredTopic filteredTopic = ImmutableFilteredTopic.getBuilder(Topic.LIVE_MESSAGES) + .withFilters(Arrays.asList("fn:filter(header:ditto-originator,'ne','some:subject')", + "fn:filter(header:ditto-origin,'ne','some-connection')")) + .build(); + + final ImmutableFilteredTopic actual = ImmutableFilteredTopic.fromString(filteredTopic.toString()); + + assertThat(actual).isEqualTo(filteredTopic); + assertThat(actual.toString()).isEqualTo(filteredTopic.toString()); + } + + @Test + public void fromStringDuplicateNamespacesParamStillThrows() { + // freezes the (out-of-scope) pre-existing behavior: only the "filter" query parameter is repeatable, + // any other duplicated parameter keeps failing like the previous Collectors.toMap-based parsing did + org.assertj.core.api.Assertions.assertThatExceptionOfType(IllegalStateException.class) + .isThrownBy(() -> ImmutableFilteredTopic.fromString( + "_/_/things/twin/events?namespaces=ns1&namespaces=ns2")) + .withMessageContaining("Duplicate key"); + } + + @Test + public void announcementTopicsDropFiltersSetViaWithFilters() { + final ImmutableFilteredTopic underTest = ImmutableFilteredTopic.getBuilder(Topic.POLICY_ANNOUNCEMENTS) + .withFilters(Arrays.asList(FILTER_EXAMPLE, "fn:filter(header:a,'exists')")) + .build(); + + assertThat(underTest.getFilters()).isEmpty(); + } + } diff --git a/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableTargetTest.java b/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableTargetTest.java index e3416b84292..030e1c8b395 100644 --- a/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableTargetTest.java +++ b/connectivity/model/src/test/java/org/eclipse/ditto/connectivity/model/ImmutableTargetTest.java @@ -14,6 +14,8 @@ import static org.assertj.core.api.Assertions.assertThat; +import java.util.Arrays; + import org.eclipse.ditto.base.model.acks.AcknowledgementLabel; import org.eclipse.ditto.base.model.auth.AuthorizationContext; import org.eclipse.ditto.base.model.auth.AuthorizationModelFactory; @@ -109,4 +111,21 @@ public void mqttFromJsonReturnsExpected() { assertThat(actual).isEqualTo(MQTT_TARGET); } + @Test + public void toJsonFromJsonRoundTripsWithRqlAndPipelineFilterTopic() { + final FilteredTopic filteredTopic = ImmutableFilteredTopic.getBuilder(Topic.TWIN_EVENTS) + .withFilters(Arrays.asList("gt(attributes/counter,42)", + "fn:filter(header:ditto-originator,'ne','some:subject')")) + .build(); + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address(ADDRESS) + .authorizationContext(AUTHORIZATION_CONTEXT) + .topics(filteredTopic) + .build(); + + final Target actual = ImmutableTarget.fromJson(target.toJson()); + + assertThat(actual).isEqualTo(target); + } + } diff --git a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActor.java b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActor.java index e2cf644aa1c..8efc4e6dd1a 100644 --- a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActor.java +++ b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActor.java @@ -24,11 +24,13 @@ import java.util.Collections; import java.util.Comparator; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; +import java.util.function.Function; import java.util.function.Predicate; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -456,7 +458,22 @@ private CompletionStage> enrichAndFilterSig // Pre-filtering already did the job return CompletableFuture.completedFuture(Collections.singletonList(outboundSignal)); } - final boolean topicWithNoFilterExists = topics.stream().anyMatch(topic -> topic.getFilter().isEmpty()); + // Partition each topic's filters exactly once: the result feeds both the aggregate boolean below and the + // per-topic filtering inside applyFilter. Topics without any filter have no map entry. + final Map partitionedFiltersByTopic = topics.stream() + .filter(topic -> !topic.getFilters().isEmpty()) + .collect(Collectors.toMap(Function.identity(), + topic -> TargetTopicFilter.partition(topic.getFilters()))); + // A topic needs no enriched thing to be decided if it has no filter at all, or if none of its filters is + // an RQL expression - a pipeline expression only ever resolves placeholders that are already known before + // enrichment (see TargetTopicFilter#matchesPipelineFilter), so applyFilter can decide such a topic even + // when signal enrichment failed and enrichedThing below ends up null. + final boolean topicWithoutThingFilterExists = topics.stream() + .anyMatch(topic -> { + @Nullable final TargetTopicFilter.PartitionedFilters partitioned = + partitionedFiltersByTopic.get(topic); + return partitioned == null || !partitioned.hasRqlExpression(); + }); final Target target = outboundSignal.getTargets().getFirst(); final DittoHeaders headers = DittoHeaders.newBuilder() @@ -495,8 +512,9 @@ private CompletionStage> enrichAndFilterSig .thenComparing(t -> t.getExtraFields().map(Object::toString).orElse("")) .thenComparing(FilteredTopic::toString)) - .filter(_ -> enrichedThing != null || topicWithNoFilterExists) - .flatMap(topic -> applyFilter(outboundSignal, enrichedThing, topic) + .filter(_ -> enrichedThing != null || topicWithoutThingFilterExists) + .flatMap(topic -> applyFilter(outboundSignal, enrichedThing, + partitionedFiltersByTopic.get(topic)) .map(signal -> enrichWithNeededExtra(signal, topic, expressionResolver, extra)) .stream()) .findFirst() @@ -825,13 +843,56 @@ private static Stream filterFailedEnrichments( } private Optional applyFilter(final OutboundSignalWithSender outboundSignal, - @Nullable final Thing thing, final FilteredTopic topic) { + @Nullable final Thing thing, + @Nullable final TargetTopicFilter.PartitionedFilters partitionedFilters) { final Signal signal = outboundSignal.getSource(); final TopicPath topicPath = DITTO_PROTOCOL_ADAPTER.toTopicPath(signal); - final Optional filter = topic.getFilter(); - if (filter.isPresent()) { + // partitionedFilters is the pre-computed partition result of the topic's filters (see enrichAndFilterSignal) + // and is null exactly when the topic has no filter at all + if (partitionedFilters != null) { + final List pipelineExpressions = partitionedFilters.getPipelineExpressions(); + if (!pipelineExpressions.isEmpty()) { + // Per the runtime failure policy, guard ONLY the pipeline evaluation: a pipeline only ever resolves + // placeholders that are already known pre-enrichment (headers, topic path, entity id, resource, + // time), so - unlike the RQL criteria below - the pipelines are evaluated first (AND, short-circuit) + // and decide a topic without RQL filter without ever needing the (possibly null, when enrichment + // failed) enriched thing. + final ExpressionResolver pipelineResolver = Resolvers.forSignal(signal, connection.getId()); + for (final String pipelineExpression : pipelineExpressions) { + final boolean pipelineMatches; + try { + pipelineMatches = + TargetTopicFilter.matchesPipelineFilter(pipelineExpression, pipelineResolver); + } catch (final DittoRuntimeException e) { + logger.withCorrelationId(signal) + .warning("Evaluating the target topic pipeline filter <{}> of connection <{}> failed " + + "with <{}>: <{}> - treating as non-match.", + pipelineExpression, connection.getId(), e.getClass().getSimpleName(), + e.getMessage()); + // an evaluation FAILURE (as opposed to an ordinary non-match, which stays silent) must be + // diagnosable by the connection owner - record it in the user-visible connection logs; + // connectionMonitorRegistry is safe to use off the actor thread (same pattern as + // logEnrichmentFailure, called from the exceptionally-stage of this future) + connectionMonitorRegistry + .forOutboundFiltered(connection, + outboundSignal.getTargets().getFirst().getOriginalAddress()) + .failure(signal, + "Evaluating the target topic pipeline filter <{0}> failed: {1} - the signal " + + "was dropped for this target topic.", + pipelineExpression, e.getMessage()); + return Optional.empty(); + } + if (!pipelineMatches) { + return Optional.empty(); + } + } + if (!partitionedFilters.hasRqlExpression()) { + // no RQL filter: decided by the pipelines alone, no thing needed + return Optional.of(outboundSignal); + } + } if (thing == null) { return Optional.empty(); } @@ -851,20 +912,26 @@ private Optional applyFilter(final OutboundSignalWithS final PlaceholderResolver timePlaceholderResolver = PlaceholderFactory .newPlaceholderResolver(TIME_PLACEHOLDER, new Object()); final DittoHeaders dittoHeaders = signal.getDittoHeaders(); - final Criteria criteria = QueryFilterCriteriaFactory.modelBased(RqlPredicateParser.getInstance(), - topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, - featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver - ).filterCriteria(filter.get(), dittoHeaders); final PlaceholderResolver thingJsonPlaceholderResolver = PlaceholderFactory .newPlaceholderResolver(THING_JSON_PLACEHOLDER, thing); - final var result = Optional.of(outboundSignal) - .filter(_ -> ThingPredicateVisitor + // at most one RQL entry after connection validation, defensively AND-combined should ever more than + // one slip through + for (final String rqlExpression : partitionedFilters.getRqlExpressions()) { + final Criteria criteria = QueryFilterCriteriaFactory.modelBased(RqlPredicateParser.getInstance(), + topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, + featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver + ).filterCriteria(rqlExpression, dittoHeaders); + final boolean rqlMatches = ThingPredicateVisitor .apply(criteria, topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver, thingJsonPlaceholderResolver) - .test(thing)); - return result; + .test(thing); + if (!rqlMatches) { + return Optional.empty(); + } + } + return Optional.of(outboundSignal); } else { // no signal enrichment: filtering is already done in SignalFilter since there is no ignored field return Optional.of(outboundSignal); diff --git a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/TargetTopicFilter.java b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/TargetTopicFilter.java new file mode 100644 index 00000000000..49d3324dd0f --- /dev/null +++ b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/TargetTopicFilter.java @@ -0,0 +1,220 @@ +/* + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.ditto.connectivity.service.messaging; + +import java.util.ArrayList; +import java.util.List; + +import org.eclipse.ditto.base.model.exceptions.DittoRuntimeException; +import org.eclipse.ditto.base.model.headers.DittoHeaders; +import org.eclipse.ditto.base.model.signals.Signal; +import org.eclipse.ditto.connectivity.model.ConnectionConfigurationInvalidException; +import org.eclipse.ditto.connectivity.model.ConnectionId; +import org.eclipse.ditto.placeholders.ExpressionResolver; +import org.eclipse.ditto.placeholders.PlaceholderFactory; + +/** + * Classifies and evaluates connection target topic filters. + *

+ * A target topic may carry up to two {@code filter} query parameters which are combined with AND semantics: at + * most one RQL expression (any value not starting with {@code fn:}, unchanged existing behavior) and at most one + * placeholder pipeline expression starting with {@code fn:} + * (e.g. {@code fn:filter(header:ditto-originator,'ne','x')}). Several {@code fn:} stages may be chained with + * {@code |} inside the pipeline parameter — a stage only runs if the previous one resolved, so chaining is AND as + * well. Both at-most-one arity rules are enforced at connection creation/update time by + * {@code ConnectionValidator}; the pipeline expression itself is validated via + * {@link #validatePipelineFilter(String, DittoHeaders)}. + *

+ * Pipeline filters should only reference headers that are stable for the signal's lifetime (such as + * {@code ditto-originator} or {@code ditto-origin}): for topics with {@code extraFields}, the pipeline is + * re-evaluated after enrichment, and internal bookkeeping headers such as {@code requested-acks} are mutated + * between the pre-enrichment gate and that re-evaluation, so filtering on them is not reliable. + */ +public final class TargetTopicFilter { + + /** + * Prefix that marks the start of a placeholder pipeline expression. + */ + private static final String FN_PREFIX = "fn:"; + + /** + * Mandatory seed prepended to every pipeline expression before evaluation. + *

+ * A pipeline that starts directly with a function invocation (e.g. {@code fn:filter(...)}) is seeded internally + * as {@link org.eclipse.ditto.placeholders.PipelineElement#unresolved()} and {@code PipelineFunctionFilter#apply} + * only ever acts {@code onResolved}. Without a resolved carrier value entering the first user-supplied function + * stage, a bare {@code fn:filter(...)} pipeline would therefore ALWAYS stay unresolved, regardless of whether the + * filter itself matches. Prepending {@code fn:default('true')|} guarantees a resolved boolean carrier value is + * fed into the first user stage, so the eventual resolved/unresolved outcome reflects the filter result rather + * than the seeding mechanics. The seed value is consumed purely as this boolean carrier and can never leak into + * the published signal. + */ + private static final String PIPELINE_SEED = "fn:default('true')|"; + + /** + * Expression resolver used for validating pipeline expressions at connection-creation/update time. + * All placeholders resolve to a dummy value; {@code ImmutableExpressionResolver} is {@code @Immutable} and + * thread-safe, so a single static instance can be shared across all validation calls. + */ + private static final ExpressionResolver VALIDATION_RESOLVER = + PlaceholderFactory.newExpressionResolverForValidation(Resolvers.getPlaceholders()); + + private TargetTopicFilter() { + throw new AssertionError(); + } + + /** + * Classifies a single target topic {@code filter} parameter value. + * + * @param filter the raw filter parameter value. + * @return {@code true} if the trimmed value starts with {@code fn:} and is therefore a placeholder pipeline + * expression, {@code false} if it is an RQL expression. + */ + public static boolean isPipelineFilter(final String filter) { + return filter.trim().startsWith(FN_PREFIX); + } + + /** + * Partitions a topic's {@code filter} parameter values into RQL and pipeline expressions, preserving order. + * Each value is trimmed and classified per {@link #isPipelineFilter(String)}. + *

+ * This method never throws and applies no structural rules: after successful connection validation each + * partition holds at most one entry, but pre-validation input may violate that - both at-most-one rules are + * enforced by {@code ConnectionValidator}, which needs the command headers for a proper error, while runtime + * callers defensively AND-evaluate whatever they get. A trimmed-empty non-{@code fn:} value stays a PRESENT + * (empty) RQL entry, so that validation keeps rejecting empty filters with an + * {@code InvalidRqlExpressionException}, exactly as an empty filter was rejected before target topic pipeline + * filters existed. + * + * @param filters the raw filter parameter values of one topic. + * @return the partitioned filter expressions. + */ + public static PartitionedFilters partition(final List filters) { + final List rqlExpressions = new ArrayList<>(1); + final List pipelineExpressions = new ArrayList<>(filters.size()); + for (final String filter : filters) { + final String trimmed = filter.trim(); + if (trimmed.startsWith(FN_PREFIX)) { + pipelineExpressions.add(trimmed); + } else { + rqlExpressions.add(trimmed); + } + } + return new PartitionedFilters(List.copyOf(rqlExpressions), List.copyOf(pipelineExpressions)); + } + + /** + * Evaluates a pipeline expression against a signal, resolving placeholders via + * {@link Resolvers#forSignal(Signal, ConnectionId)}. + * + * @param pipelineExpression the pipeline expression (without the mandatory seed). + * @param signal the signal the filter is evaluated against. + * @param connectionId the ID of the connection evaluating the filter. + * @return {@code true} if the pipeline resolves (i.e. the filter matches and the target topic should be + * published), {@code false} if it stays unresolved (i.e. the topic should be suppressed). + * @throws DittoRuntimeException if the pipeline expression is malformed or cannot be evaluated. Callers at + * runtime sites are responsible for catching this per the runtime failure policy and treating it as a + * non-match. + */ + public static boolean matchesPipelineFilter(final String pipelineExpression, final Signal signal, + final ConnectionId connectionId) { + return matchesPipelineFilter(pipelineExpression, Resolvers.forSignal(signal, connectionId)); + } + + /** + * Evaluates a pipeline expression using the given expression resolver. + * + * @param pipelineExpression the pipeline expression (without the mandatory seed). + * @param resolver the expression resolver to resolve placeholders in the pipeline expression with. + * @return {@code true} if the pipeline resolves (i.e. the filter matches and the target topic should be + * published), {@code false} if it stays unresolved (i.e. the topic should be suppressed). + * @throws DittoRuntimeException if the pipeline expression is malformed or cannot be evaluated. Callers at + * runtime sites are responsible for catching this per the runtime failure policy and treating it as a + * non-match. + */ + public static boolean matchesPipelineFilter(final String pipelineExpression, final ExpressionResolver resolver) { + return resolver.resolveAsPipelineElement(PIPELINE_SEED + pipelineExpression).findFirst().isPresent(); + } + + /** + * Validates a pipeline expression at connection-creation/update time, i.e. strictly: any placeholder/pipeline + * function error is rejected. Several {@code fn:} stages may be chained with {@code |} inside the one pipeline + * filter parameter; the resolver's pipeline grammar enforces the structure (quote-aware stage splitting, every + * stage a function invocation, at most 10 {@code fn:} stages). + * + * @param pipelineExpression the pipeline expression (without the mandatory seed) to validate. + * @param dittoHeaders the headers of the command which triggered the validation, stamped onto the thrown + * exception for correlation. + * @throws ConnectionConfigurationInvalidException if the pipeline expression is invalid, e.g. because it + * references an unknown placeholder function, has an invalid function signature, exceeds the maximum number + * of pipeline stages, or references an unresolvable placeholder. + */ + public static void validatePipelineFilter(final String pipelineExpression, final DittoHeaders dittoHeaders) { + try { + VALIDATION_RESOLVER.resolveAsPipelineElement(PIPELINE_SEED + pipelineExpression); + } catch (final DittoRuntimeException e) { + throw ConnectionConfigurationInvalidException + .newBuilder("The target topic pipeline filter expression '" + pipelineExpression + + "' is invalid: " + e.getMessage()) + .description(e.getDescription() + .orElse("Check the spelling and syntax of the pipeline expression.") + ) + .cause(e) + .dittoHeaders(dittoHeaders) + .build(); + } + } + + /** + * Immutable result of {@link TargetTopicFilter#partition(List)}: a topic's filter parameter values, partitioned + * into RQL and pipeline expressions with their relative order preserved. + *

+ * An empty string is a PRESENT (if empty) RQL entry, not an absent one - see + * {@link TargetTopicFilter#partition(List)} for why this matters for validation. + */ + public static final class PartitionedFilters { + + private final List rqlExpressions; + private final List pipelineExpressions; + + private PartitionedFilters(final List rqlExpressions, final List pipelineExpressions) { + this.rqlExpressions = rqlExpressions; + this.pipelineExpressions = pipelineExpressions; + } + + /** + * @return the RQL expressions among the topic's filters - at most one entry after successful connection + * validation, but possibly more for not (yet) validated input. + */ + public List getRqlExpressions() { + return rqlExpressions; + } + + /** + * @return the pipeline expressions (each without the mandatory seed) among the topic's filters - at most + * one entry after successful connection validation, but possibly more for not (yet) validated input. + */ + public List getPipelineExpressions() { + return pipelineExpressions; + } + + /** + * @return whether at least one of the topic's filters is an RQL expression. + */ + public boolean hasRqlExpression() { + return !rqlExpressions.isEmpty(); + } + + } + +} diff --git a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilter.java b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilter.java index 5bf7a302bfb..06314b875f8 100644 --- a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilter.java +++ b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilter.java @@ -26,6 +26,7 @@ import org.eclipse.ditto.base.model.auth.AuthorizationContext; import org.eclipse.ditto.base.model.entity.id.EntityId; import org.eclipse.ditto.base.model.entity.id.WithEntityId; +import org.eclipse.ditto.base.model.exceptions.DittoRuntimeException; import org.eclipse.ditto.base.model.headers.DittoHeaders; import org.eclipse.ditto.base.model.namespaces.NamespaceReader; import org.eclipse.ditto.base.model.signals.Signal; @@ -34,19 +35,25 @@ import org.eclipse.ditto.base.model.signals.commands.CommandResponse; import org.eclipse.ditto.base.model.signals.events.Event; import org.eclipse.ditto.connectivity.model.Connection; +import org.eclipse.ditto.connectivity.model.ConnectionId; import org.eclipse.ditto.connectivity.model.FilteredTopic; import org.eclipse.ditto.connectivity.model.Target; import org.eclipse.ditto.connectivity.model.Topic; import org.eclipse.ditto.connectivity.model.signals.announcements.ConnectivityAnnouncement; +import org.eclipse.ditto.connectivity.service.messaging.Resolvers; +import org.eclipse.ditto.connectivity.service.messaging.TargetTopicFilter; import org.eclipse.ditto.connectivity.service.messaging.monitoring.ConnectionMonitor; import org.eclipse.ditto.connectivity.service.messaging.monitoring.ConnectionMonitorRegistry; import org.eclipse.ditto.edge.service.placeholders.EntityIdPlaceholder; import org.eclipse.ditto.edge.service.placeholders.FeaturePlaceholder; import org.eclipse.ditto.edge.service.placeholders.ThingPlaceholder; +import org.eclipse.ditto.internal.utils.pekko.logging.DittoLogger; +import org.eclipse.ditto.internal.utils.pekko.logging.DittoLoggerFactory; import org.eclipse.ditto.json.JsonFieldSelector; import org.eclipse.ditto.json.JsonPointer; import org.eclipse.ditto.messages.model.signals.commands.MessageCommand; import org.eclipse.ditto.messages.model.signals.commands.MessageCommandResponse; +import org.eclipse.ditto.placeholders.ExpressionResolver; import org.eclipse.ditto.placeholders.PlaceholderFactory; import org.eclipse.ditto.placeholders.PlaceholderResolver; import org.eclipse.ditto.placeholders.TimePlaceholder; @@ -72,6 +79,8 @@ */ public final class SignalFilter { + private static final DittoLogger LOGGER = DittoLoggerFactory.getLogger(SignalFilter.class); + private static final DittoProtocolAdapter DITTO_PROTOCOL_ADAPTER = DittoProtocolAdapter.newInstance(); private static final TopicPathPlaceholder TOPIC_PATH_PLACEHOLDER = TopicPathPlaceholder.getInstance(); private static final EntityIdPlaceholder ENTITY_ID_PLACEHOLDER = EntityIdPlaceholder.getInstance(); @@ -104,14 +113,26 @@ public static SignalFilter of(final Connection connection, /** * Filters the passed {@code signal} by extracting those {@link Target}s which should receive the signal. * Fields are ignored if they occur as "extra targets" to be evaluated later after signal enrichment. + *

+ * A target topic may carry up to two {@code filter} parameters, combined with AND semantics: at most one RQL + * expression (unchanged, existing behavior) plus at most one placeholder pipeline expression ({@code fn:...}, + * possibly chaining several stages with {@code |} - see + * {@link org.eclipse.ditto.connectivity.service.messaging.TargetTopicFilter}); the loop below defensively + * AND-evaluates every pipeline entry it finds, even for never-validated topics carrying more. The + * pipeline filters are evaluated first, before enrichment, as a deterministic hard gate; per the runtime + * failure policy, a {@link org.eclipse.ditto.base.model.exceptions.DittoRuntimeException} thrown while + * evaluating one of them is caught, logged as a warning plus a failure entry in the user-visible connection + * logs, and treated as a non-match rather than propagated. The RQL filter - if present - keeps its existing + * (unguarded) behavior. * * @param signal the signal to filter / determine the {@link org.eclipse.ditto.connectivity.model.Target}s for * @return the determined Targets for the passed in {@code signal} - * @throws org.eclipse.ditto.base.model.exceptions.InvalidRqlExpressionException if the optional filter string of a - * Target cannot be mapped to a valid criterion + * @throws org.eclipse.ditto.base.model.exceptions.InvalidRqlExpressionException if the optional RQL filter of a + * Target's topic cannot be mapped to a valid criterion */ @SuppressWarnings("squid:S3864") public List filter(final Signal signal) { + final ConnectionId connectionId = connection.getId(); return connection.getTargets().stream() .filter(t -> isTargetAuthorized(t, signal)) // this is cheaper, so check this first .filter(t -> isTargetSubscribedForTopicGenerally(t, signal)) @@ -119,7 +140,7 @@ public List filter(final Signal signal) { .peek(authorizedTarget -> connectionMonitorRegistry.forOutboundDispatched(connection, authorizedTarget.getAddress()) .success(signal)) - .filter(t -> isTargetSubscribedForTopicWithFiltering(t, signal)) + .filter(t -> isTargetSubscribedForTopicWithFiltering(t, signal, connectionId)) // count authorized + filtered targets .peek(filteredTarget -> connectionMonitorRegistry.forOutboundFiltered(connection, filteredTarget.getAddress()) @@ -143,11 +164,12 @@ private static boolean isTargetSubscribedForTopicGenerally(final Target target, .anyMatch(applyTopicFilter(signal)); } - private static boolean isTargetSubscribedForTopicWithFiltering(final Target target, final Signal signal) { + private boolean isTargetSubscribedForTopicWithFiltering(final Target target, final Signal signal, + final ConnectionId connectionId) { return target.getTopics().stream() .filter(applyTopicFilter(signal)) .filter(applyNamespaceFilter(signal)) - .anyMatch(filteredTopic -> matchesFilterBeforeEnrichment(filteredTopic, signal)); + .anyMatch(filteredTopic -> matchesFilterBeforeEnrichment(filteredTopic, target, signal, connectionId)); } private static Predicate applyTopicFilter(final Signal signal) { @@ -164,47 +186,93 @@ private static String namespaceFromId(final WithEntityId withEntityId) { return NamespaceReader.fromEntityId(withEntityId.getEntityId()).orElse(null); } - private static boolean matchesFilterBeforeEnrichment(final FilteredTopic filteredTopic, final Signal signal) { - final Optional filterOptional = filteredTopic.getFilter(); - if (filterOptional.isPresent()) { - // match filter ignoring "extraFields" + private boolean matchesFilterBeforeEnrichment(final FilteredTopic filteredTopic, final Target target, + final Signal signal, final ConnectionId connectionId) { + final List filters = filteredTopic.getFilters(); + if (filters.isEmpty()) { + return true; + } + final TargetTopicFilter.PartitionedFilters partitioned = TargetTopicFilter.partition(filters); + final List pipelineExpressions = partitioned.getPipelineExpressions(); + if (!pipelineExpressions.isEmpty()) { + // The pipeline filters are a deterministic hard gate evaluated BEFORE enrichment: unlike the RQL + // criteria below - which need a thing snapshot reconstructed from the event and are therefore + // only meaningfully evaluable for ThingEvents - a pipeline only ever resolves placeholders + // that are already fully known pre-enrichment (headers, topic path, entity id, resource, time), + // for ANY filterable signal type (twin/live events, live commands, live messages alike). Its + // match/non-match outcome can therefore never change once/if enrichment happens, so a single + // non-match can short-circuit the whole target right here, and - for a topic without an RQL filter - + // an all-match makes the target immediately eligible without ever touching the RQL path below. + final ExpressionResolver expressionResolver = Resolvers.forSignal(signal, connectionId); + for (final String pipelineExpression : pipelineExpressions) { + final boolean pipelineMatches; + try { + pipelineMatches = + TargetTopicFilter.matchesPipelineFilter(pipelineExpression, expressionResolver); + } catch (final DittoRuntimeException e) { + LOGGER.withCorrelationId(signal) + .warn("Evaluating the target topic pipeline filter <{}> of connection <{}> failed with " + + "<{}>: <{}> - treating as non-match.", + pipelineExpression, connectionId, e.getClass().getSimpleName(), e.getMessage()); + // an evaluation FAILURE (as opposed to an ordinary non-match, which stays silent) must be + // diagnosable by the connection owner - record it in the user-visible connection logs + connectionMonitorRegistry.forOutboundFiltered(connection, target.getAddress()) + .failure(signal, + "Evaluating the target topic pipeline filter <{0}> failed: {1} - the signal " + + "was dropped for this target topic.", + pipelineExpression, e.getMessage()); + return false; + } + if (!pipelineMatches) { + return false; + } + } + } + if (!partitioned.hasRqlExpression()) { + return true; + } - final TopicPath topicPath = DITTO_PROTOCOL_ADAPTER.toTopicPath(signal); - final PlaceholderResolver topicPathPlaceholderResolver = - PlaceholderFactory.newPlaceholderResolver(TOPIC_PATH_PLACEHOLDER, topicPath); - final PlaceholderResolver entityIdPlaceholderResolver = PlaceholderFactory - .newPlaceholderResolver(ENTITY_ID_PLACEHOLDER, - (signal instanceof WithEntityId withEntityId) ? withEntityId.getEntityId() : null); - final PlaceholderResolver thingPlaceholderResolver = PlaceholderFactory - .newPlaceholderResolver(THING_PLACEHOLDER, - (signal instanceof WithEntityId withEntityId) ? withEntityId.getEntityId() : null); - final PlaceholderResolver> featurePlaceholderResolver = PlaceholderFactory - .newPlaceholderResolver(FEATURE_PLACEHOLDER, signal); - final PlaceholderResolver resourcePlaceholderResolver = PlaceholderFactory - .newPlaceholderResolver(RESOURCE_PLACEHOLDER, signal); - final PlaceholderResolver timePlaceholderResolver = PlaceholderFactory - .newPlaceholderResolver(TIME_PLACEHOLDER, new Object()); - final Criteria criteria = parseCriteria(filterOptional.get(), signal.getDittoHeaders(), + // match RQL filter(s) ignoring "extraFields" - at most one entry after connection validation, defensively + // AND-combined should ever more than one slip through + final TopicPath topicPath = DITTO_PROTOCOL_ADAPTER.toTopicPath(signal); + final PlaceholderResolver topicPathPlaceholderResolver = + PlaceholderFactory.newPlaceholderResolver(TOPIC_PATH_PLACEHOLDER, topicPath); + final PlaceholderResolver entityIdPlaceholderResolver = PlaceholderFactory + .newPlaceholderResolver(ENTITY_ID_PLACEHOLDER, + (signal instanceof WithEntityId withEntityId) ? withEntityId.getEntityId() : null); + final PlaceholderResolver thingPlaceholderResolver = PlaceholderFactory + .newPlaceholderResolver(THING_PLACEHOLDER, + (signal instanceof WithEntityId withEntityId) ? withEntityId.getEntityId() : null); + final PlaceholderResolver> featurePlaceholderResolver = PlaceholderFactory + .newPlaceholderResolver(FEATURE_PLACEHOLDER, signal); + final PlaceholderResolver resourcePlaceholderResolver = PlaceholderFactory + .newPlaceholderResolver(RESOURCE_PLACEHOLDER, signal); + final PlaceholderResolver timePlaceholderResolver = PlaceholderFactory + .newPlaceholderResolver(TIME_PLACEHOLDER, new Object()); + final Set extraFields = filteredTopic.getExtraFields() + .map(JsonFieldSelector::getPointers) + .orElse(Collections.emptySet()); + final Thing thingToMatch; + if (signal instanceof ThingEvent) { + final Optional thingFromEvent = ThingEventToThingConverter.thingEventToThing((ThingEvent) signal); + if (thingFromEvent.isEmpty()) { + return false; + } + thingToMatch = thingFromEvent.get(); + } else { + thingToMatch = Thing.newBuilder().build(); + } + for (final String filter : partitioned.getRqlExpressions()) { + final Criteria criteria = parseCriteria(filter, signal.getDittoHeaders(), topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver); - final Set extraFields = filteredTopic.getExtraFields() - .map(JsonFieldSelector::getPointers) - .orElse(Collections.emptySet()); - if (signal instanceof ThingEvent) { - return ThingEventToThingConverter.thingEventToThing((ThingEvent) signal) - .filter(thing -> Thing3ValuePredicateVisitor.couldBeTrue(criteria, extraFields, thing, - topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, - featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver)) - .isPresent(); - } else { - final Thing emptyThing = Thing.newBuilder().build(); - return Thing3ValuePredicateVisitor.couldBeTrue(criteria, extraFields, emptyThing, - topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, - featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver); + if (!Thing3ValuePredicateVisitor.couldBeTrue(criteria, extraFields, thingToMatch, + topicPathPlaceholderResolver, entityIdPlaceholderResolver, thingPlaceholderResolver, + featurePlaceholderResolver, resourcePlaceholderResolver, timePlaceholderResolver)) { + return false; } - } else { - return true; } + return true; } /** diff --git a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidator.java b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidator.java index 8ec0aaf19d2..3e057b6300c 100644 --- a/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidator.java +++ b/connectivity/service/src/main/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidator.java @@ -46,6 +46,7 @@ import org.eclipse.ditto.connectivity.service.config.ConnectionConfig; import org.eclipse.ditto.connectivity.service.config.ConnectivityConfig; import org.eclipse.ditto.connectivity.service.config.mapping.MapperLimitsConfig; +import org.eclipse.ditto.connectivity.service.messaging.TargetTopicFilter; import org.eclipse.ditto.connectivity.service.messaging.internal.ssl.SSLContextCreator; import org.eclipse.ditto.connectivity.service.messaging.monitoring.logs.ConnectionLogger; import org.eclipse.ditto.connectivity.service.placeholders.ConnectivityPlaceholders; @@ -329,10 +330,44 @@ private void validateSourceAndTargetAddressesAreNonempty(final Connection connec final String location = String.format("Targets of connection <%s>", connection.getId()); throw emptyAddressesError(location, dittoHeaders); } - target.getTopics().forEach(topic -> topic.getFilter().ifPresent(filter -> - // will throw an InvalidRqlExpressionException if the RQL expression was not valid: - queryFilterCriteriaFactory.filterCriteria(filter, dittoHeaders) - )); + target.getTopics().forEach(topic -> { + final List filters = topic.getFilters(); + if (!filters.isEmpty()) { + final TargetTopicFilter.PartitionedFilters partitionedFilters = + TargetTopicFilter.partition(filters); + if (partitionedFilters.getRqlExpressions().size() > 1) { + throw ConnectionConfigurationInvalidException + .newBuilder("The topic '" + topic + "' of the target with address '" + + target.getAddress() + "' declares " + + partitionedFilters.getRqlExpressions().size() + " RQL 'filter' parameters " + + "- at most one RQL filter is allowed per topic.") + .description("Combine several RQL conditions into a single RQL expression using " + + "'and(...)'. Besides the RQL filter, a topic may only carry one more " + + "'filter' parameter holding a placeholder pipeline expression starting " + + "with 'fn:'.") + .dittoHeaders(dittoHeaders) + .build(); + } + if (partitionedFilters.getPipelineExpressions().size() > 1) { + throw ConnectionConfigurationInvalidException + .newBuilder("The topic '" + topic + "' of the target with address '" + + target.getAddress() + "' declares " + + partitionedFilters.getPipelineExpressions().size() + + " pipeline 'filter' parameters - at most one pipeline filter is allowed " + + "per topic.") + .description("Combine several pipeline conditions by chaining 'fn:' stages with " + + "'|' inside the single pipeline 'filter' parameter instead, e.g. " + + "'?filter=fn:filter(...)|fn:filter(...)'.") + .dittoHeaders(dittoHeaders) + .build(); + } + partitionedFilters.getPipelineExpressions().forEach(pipeline -> + TargetTopicFilter.validatePipelineFilter(pipeline, dittoHeaders)); + partitionedFilters.getRqlExpressions().forEach(rql -> + // will throw an InvalidRqlExpressionException if the RQL expression was not valid: + queryFilterCriteriaFactory.filterCriteria(rql, dittoHeaders)); + } + }); }); } diff --git a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActorTest.java b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActorTest.java index b5a6b36d315..00e2d9c8a31 100644 --- a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActorTest.java +++ b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/OutboundMappingProcessorActorTest.java @@ -726,6 +726,295 @@ public void crossStreamingTypeCatchAllDoesNotLeakIntoOtherStreamingType() { }}; } + // --- target topic pipeline (fn:) filter tests (Task 4) --- + + @Test + public void mixedTopicsTargetPurePipelineNonMatchDropsEntireTarget() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // Target with an RQL+extraFields topic that cannot match this signal (topic4 requires + // resource:path == /features/feature4), plus a pure-pipeline topic on "ditto-originator". + // Neither topic matches, so the whole target is dropped. + final FilteredTopic pipelineTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:filter(header:ditto-originator,'ne','x')") + .build(); + final List targets = List.of(createTestTargetMultiTopics(Set.of(topic4(), pipelineTopic))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("unrelatedFeature").build(), + List.of("source1", "multipleExtraFields"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + // THEN: sender receives a weak acknowledgement for the dropped target - no signal is published + final Acknowledgements acks = expectMsgClass(Acknowledgements.class); + final List ackLabels = acks.getSuccessfulAcknowledgements() + .stream() + .map(ack -> ack.getLabel().toString()) + .toList(); + assertThat(ackLabels).containsExactlyInAnyOrder("source1", "multipleExtraFields"); + acks.forEach(ack -> assertThat(ack.isWeak()).describedAs("Expect weak ack, got: " + ack).isTrue()); + }}; + } + + @Test + public void mixedTopicsTargetPurePipelineMatchPublishesWithoutExtraFieldsAndHeadersUnchanged() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // Same mixed target as above, but the pipeline now matches: the target must publish via the + // pure-pipeline topic (no thing needed) without any extra fields, and without leaking anything + // (in particular not the internal pipeline seed, see TargetTopicFilter#PIPELINE_SEED) onto the + // published signal's headers. A second, unfiltered "control" target on the very same signal + // establishes what the mapping/dispatch pipeline normally does to headers (e.g. stripping + // internal-only ones, adding "content-type") so the comparison isolates exactly what the pipeline + // *filter evaluation* itself changed, rather than unrelated header bookkeeping. + final FilteredTopic pipelineTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:filter(header:ditto-originator,'eq','x')") + .build(); + final Target pipelineTarget = createTestTargetMultiTopics(Set.of(topic4(), pipelineTopic)); + final Target controlTarget = ConnectivityModelFactory.newTargetBuilder(pipelineTarget) + .address("pipelineControlTarget") + .topics(Set.of(topic0())) + .issuedAcknowledgementLabel(AcknowledgementLabel.of("pipelineControl")) + .build(); + final List targets = List.of(pipelineTarget, controlTarget); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("unrelatedFeature").build(), + List.of("source1", "multipleExtraFields", "pipelineControl"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + final BaseClientActor.PublishMappedMessage publish = + clientActorProbe.expectMsgClass(BaseClientActor.PublishMappedMessage.class); + + assertThat(publish.getOutboundSignal().getMappedOutboundSignals()).hasSize(2); + assertThat(publish.getOutboundSignal().first().getTargets()).contains(pipelineTarget); + + // Verify no extra fields: the winning topic is the pure-pipeline one, which carries none + assertThat(publish.getOutboundSignal().first().getAdaptable().getPayload().getExtra()) + .as("Outbound signal must not contain any extra field: the pure-pipeline topic won selection.") + .isEmpty(); + + // Verify no leak: the "ditto-originator" header survives untouched, and the header key set is + // identical to the control target's - i.e. the pipeline evaluation (in particular the internal + // boolean seed) added/removed nothing beyond what plain unfiltered mapping already does. + final DittoHeaders pipelineHeaders = publish.getOutboundSignal().first().getAdaptable().getDittoHeaders(); + final DittoHeaders controlHeaders = + publish.getOutboundSignal().getMappedOutboundSignals().get(1).getAdaptable().getDittoHeaders(); + assertThat(pipelineHeaders.get("ditto-originator")).isEqualTo("x"); + assertThat(pipelineHeaders.keySet()) + .as("Pipeline evaluation must not add/remove any header compared to an unfiltered publish of " + + "the same signal") + .isEqualTo(controlHeaders.keySet()); + }}; + } + + @Test + public void mixedTopicsTargetEnrichmentSucceedsPipelineMatchesPublishesWithRqlTopicExtraFields() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // T-C1 scenario 1: a single target carries both an RQL+extraFields topic (topic4) and a + // pure-pipeline topic. Enrichment succeeds and both independently match; the extraFields-first + // sort (see enrichAndFilterSignal) must make topic4 win the extraFields selection. + final FilteredTopic pipelineTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:filter(header:ditto-originator,'eq','x')") + .build(); + final List targets = List.of(createTestTargetMultiTopics(Set.of(topic4(), pipelineTopic))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("feature4") + .build().setProperties(FeatureProperties.newBuilder().set("size", "large").build()), + List.of("multipleExtraFields"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + final BaseClientActor.PublishMappedMessage publish = + clientActorProbe.expectMsgClass(BaseClientActor.PublishMappedMessage.class); + + assertThat(publish.getOutboundSignal().first().getTargets()).contains(targets.getFirst()); + assertThat(publish.getOutboundSignal().first().getAdaptable().getPayload().getExtra()) + .isPresent() + .hasValueSatisfying(extra -> + assertThat(extra.getValue(JsonPointer.of("definition"))) + .as("Outbound signal does not contain the requested extra fields from topic4.") + .isPresent()); + }}; + } + + @Test + public void mixedTopicsTargetEnrichmentSucceedsPipelineNonMatchRqlMatchesPublishesWithExtraFields() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // T-C1 scenario 2: same mixed target, but the pipeline does NOT match this time (originator "x" != + // required "y"). topic4's RQL matches independently on feature4, so the target must still publish + // with topic4's extraFields - a non-matching pipeline topic must not sink the whole target nor win + // the extraFields selection (it must return empty from applyFilter). + final FilteredTopic pipelineTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:filter(header:ditto-originator,'eq','y')") + .build(); + final List targets = List.of(createTestTargetMultiTopics(Set.of(topic4(), pipelineTopic))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("feature4") + .build().setProperties(FeatureProperties.newBuilder().set("size", "large").build()), + List.of("multipleExtraFields"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + final BaseClientActor.PublishMappedMessage publish = + clientActorProbe.expectMsgClass(BaseClientActor.PublishMappedMessage.class); + + assertThat(publish.getOutboundSignal().first().getTargets()).contains(targets.getFirst()); + assertThat(publish.getOutboundSignal().first().getAdaptable().getPayload().getExtra()) + .isPresent() + .hasValueSatisfying(extra -> + assertThat(extra.getValue(JsonPointer.of("definition"))) + .as("RQL topic must win and enrich the outbound signal even though the " + + "sibling pipeline topic did not match.") + .isPresent()); + }}; + } + + @Test + public void rqlAndPipelineFilterParamsEvaluateRqlAgainstEnrichedThing() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // A single topic with an RQL filter param, a pipeline filter param (AND semantics) and extraFields. + // The RQL param references "features/featureA", which is absent from the signal-derived (raw) thing + // and only appears once partialRetrieveAndResponse's extra data is merged in - so a match proves the + // RQL param (fed via TargetTopicFilter.PartitionedFilters#getRqlExpressions) is evaluated against + // the enriched thing. + final FilteredTopic combinedTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of("exists(features/featureA)", + "fn:filter(header:ditto-originator,'eq','x')")) + .withExtraFields(ThingFieldSelector.fromString("definition")) + .build(); + final List targets = List.of(createTestTargetMultiTopics(Set.of(combinedTopic))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("unrelatedFeature").build(), + List.of("multipleExtraFields"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + final BaseClientActor.PublishMappedMessage publish = + clientActorProbe.expectMsgClass(BaseClientActor.PublishMappedMessage.class); + + assertThat(publish.getOutboundSignal().first().getTargets()).contains(targets.getFirst()); + assertThat(publish.getOutboundSignal().first().getAdaptable().getPayload().getExtra()) + .isPresent() + .hasValueSatisfying(extra -> assertThat(extra.getValue(JsonPointer.of("definition"))) + .isPresent()); + }}; + } + + @Test + public void chainedPipelineStagesWithExtraFieldsAllMatchPublishesEnriched() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // A single topic with one pipeline filter param chaining TWO fn: stages (AND semantics) and + // extraFields: both stages match, so the signal is published with the topic's extra fields after + // the post-enrichment re-evaluation. + final FilteredTopic pipelineTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of("fn:filter(header:ditto-originator,'eq','x')" + + "|fn:filter(header:ditto-originator,'ne','y')")) + .withExtraFields(ThingFieldSelector.fromString("definition")) + .build(); + final List targets = List.of(createTestTargetMultiTopics(Set.of(pipelineTopic))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("unrelatedFeature").build(), + List.of("multipleExtraFields"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + final BaseClientActor.PublishMappedMessage publish = + clientActorProbe.expectMsgClass(BaseClientActor.PublishMappedMessage.class); + + assertThat(publish.getOutboundSignal().first().getTargets()).contains(targets.getFirst()); + assertThat(publish.getOutboundSignal().first().getAdaptable().getPayload().getExtra()) + .isPresent() + .hasValueSatisfying(extra -> assertThat(extra.getValue(JsonPointer.of("definition"))) + .isPresent()); + }}; + } + + @Test + public void chainedPipelineStagesSecondStageNonMatchDropsTarget() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // Same topic shape as above, but the SECOND chained stage does not match ("ne 'x'" with originator + // "x") - AND semantics must drop the whole target even though the first stage matches. + final FilteredTopic pipelineTopic = ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of("fn:filter(header:ditto-originator,'eq','x')" + + "|fn:filter(header:ditto-originator,'ne','x')")) + .withExtraFields(ThingFieldSelector.fromString("definition")) + .build(); + final List targets = List.of(createTestTargetMultiTopics(Set.of(pipelineTopic))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = withOriginator(outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("unrelatedFeature").build(), + List.of("source1", "multipleExtraFields"), targets, getRef()), "x"); + + underTest.tell(outboundSignal, getRef()); + partialRetrieveAndResponse(); + + // THEN: sender receives weak acknowledgements for the dropped target - no signal is published + final Acknowledgements acks = expectMsgClass(Acknowledgements.class); + final List ackLabels = acks.getSuccessfulAcknowledgements() + .stream() + .map(ack -> ack.getLabel().toString()) + .toList(); + assertThat(ackLabels).containsExactlyInAnyOrder("source1", "multipleExtraFields"); + acks.forEach(ack -> assertThat(ack.isWeak()).describedAs("Expect weak ack, got: " + ack).isTrue()); + }}; + } + + @Test + public void rqlExtraFieldsOnlyTargetStillDroppedWhenEnrichmentFails() { + new TestKit(actorSystemResource.getActorSystem()) {{ + // Regression guard for T-C1 scenario 4: a target with only an RQL+extraFields topic (no pipeline + // topic) must still be dropped when signal enrichment genuinely fails - the widened :459/:498 + // predicate must not change this pre-existing behavior. Enrichment failure is simulated exactly + // like the existing eventsWithFailedEnrichmentIssueFailedAcks test: never reply to the + // RetrieveThing, so the underlying ask-with-retry ultimately fails. That failure is intercepted by + // OutboundMappingProcessorActor's own recovery path (recoverFromEnrichmentError) before + // enrichAndFilterSignal's topics-selection code ever runs, so it proves the observable + // "still dropped" contract end-to-end. + final List targets = List.of(createTestTargetMultiTopics(Set.of(topic4()))); + final Connection connection = CONNECTION.toBuilder().setTargets(targets).build(); + final ActorRef underTest = getTestActorRef(connection); + + final OutboundSignal outboundSignal = outboundFeatureTwinEvent(THING, + Feature.newBuilder().withId("feature4").build(), + List.of("source1", "multipleExtraFields"), targets, getRef()); + underTest.tell(outboundSignal, getRef()); + proxyActorProbe.expectMsgClass(RetrieveThing.class); + // no reply: the retrieval ask fails/times out, simulating an enrichment failure + + final Acknowledgements acks = expectMsgClass(Duration.ofSeconds(5), Acknowledgements.class); + final List failedAckLabels = acks.getFailedAcknowledgements() + .stream() + .map(ack -> ack.getLabel().toString()) + .toList(); + assertThat(failedAckLabels).containsExactly("multipleExtraFields"); + }}; + } + private void partialRetrieveAndResponse() { // Expect enrichment request for all fields in all topics final RetrieveThing retrieveEnrichedThing = proxyActorProbe.expectMsgClass(RetrieveThing.class); @@ -810,6 +1099,18 @@ private static OutboundSignal outboundLiveEvent(final Attributes attributes, fin outboundTwinEvent.getTargets()); } + /** + * Returns a copy of {@code outboundSignal} whose source signal carries an explicit {@code ditto-originator} + * header, as required by any target topic pipeline filter relying on {@code header:ditto-originator}. + */ + private static OutboundSignal withOriginator(final OutboundSignal outboundSignal, final String originator) { + return OutboundSignalFactory.newOutboundSignal( + outboundSignal.getSource().setDittoHeaders(outboundSignal.getSource().getDittoHeaders().toBuilder() + .putHeader("ditto-originator", originator) + .build()), + outboundSignal.getTargets()); + } + private static Connection createTestConnection() { final ConnectionType type = ConnectionType.MQTT_5; final ConnectivityStatus status = ConnectivityStatus.OPEN; diff --git a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/TargetTopicFilterTest.java b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/TargetTopicFilterTest.java new file mode 100644 index 00000000000..d39325f9d29 --- /dev/null +++ b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/TargetTopicFilterTest.java @@ -0,0 +1,363 @@ +/* + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.ditto.connectivity.service.messaging; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; +import static org.assertj.core.api.Assertions.assertThatNoException; +import static org.assertj.core.api.Assertions.catchThrowable; + +import java.time.Instant; +import java.util.Collections; +import java.util.List; +import java.util.Map; + +import org.eclipse.ditto.base.model.headers.DittoHeaders; +import org.eclipse.ditto.base.model.signals.Signal; +import org.eclipse.ditto.connectivity.model.ConnectionConfigurationInvalidException; +import org.eclipse.ditto.connectivity.model.ConnectionId; +import org.eclipse.ditto.connectivity.service.messaging.TargetTopicFilter.PartitionedFilters; +import org.eclipse.ditto.json.JsonPointer; +import org.eclipse.ditto.json.JsonValue; +import org.eclipse.ditto.things.model.Thing; +import org.eclipse.ditto.things.model.ThingId; +import org.eclipse.ditto.things.model.signals.events.ThingModified; +import org.junit.Test; + +/** + * Tests {@link TargetTopicFilter}. + */ +public final class TargetTopicFilterTest { + + private static final ConnectionId CONNECTION_ID = ConnectionId.generateRandom(); + private static final ThingId THING_ID = ThingId.of("foo:bar13"); + + // ===== isPipelineFilter(): classification ===== + + @Test + public void isPipelineFilterClassifiesByFnPrefix() { + assertThat(TargetTopicFilter.isPipelineFilter("fn:filter(header:a,'exists')")).isTrue(); + assertThat(TargetTopicFilter.isPipelineFilter(" fn:filter(header:a,'exists')")).isTrue(); + assertThat(TargetTopicFilter.isPipelineFilter("gt(attributes/x,5)")).isFalse(); + // an "fn:" substring anywhere but the (trimmed) start does not make a pipeline filter + assertThat(TargetTopicFilter.isPipelineFilter("like(attributes/a,'*|fn:x*')")).isFalse(); + assertThat(TargetTopicFilter.isPipelineFilter("")).isFalse(); + } + + // ===== partition(): pure forms ===== + + @Test + public void partitionPurePipelineFilter() { + final PartitionedFilters partitioned = + TargetTopicFilter.partition(List.of("fn:filter(header:ditto-originator,'eq','x')")); + + assertThat(partitioned.getRqlExpressions()).isEmpty(); + assertThat(partitioned.hasRqlExpression()).isFalse(); + assertThat(partitioned.getPipelineExpressions()) + .containsExactly("fn:filter(header:ditto-originator,'eq','x')"); + } + + @Test + public void partitionPureRqlFilters() { + assertPureRql("and(eq(attributes/a,1),eq(attributes/b,2))"); + assertPureRql("eq(attributes/a,1)"); + assertPureRql("exists(attributes/a)"); + } + + private static void assertPureRql(final String rql) { + final PartitionedFilters partitioned = TargetTopicFilter.partition(List.of(rql)); + + assertThat(partitioned.getRqlExpressions()).containsExactly(rql); + assertThat(partitioned.getPipelineExpressions()).isEmpty(); + } + + // ===== partition(): mixed parameter lists ===== + + @Test + public void partitionClassifiesMixedParamsPreservingOrder() { + final PartitionedFilters partitioned = TargetTopicFilter.partition(List.of( + "gt(attributes/x,5)", + "fn:filter(header:a,'exists')", + "fn:filter(header:b,'ne','x')")); + + assertThat(partitioned.getRqlExpressions()).containsExactly("gt(attributes/x,5)"); + assertThat(partitioned.hasRqlExpression()).isTrue(); + assertThat(partitioned.getPipelineExpressions()) + .containsExactly("fn:filter(header:a,'exists')", "fn:filter(header:b,'ne','x')"); + } + + @Test + public void partitionCollectsMultipleRqlParamsWithoutThrowing() { + // partition() applies no structural rules - the at-most-one-RQL rule is enforced by ConnectionValidator + final PartitionedFilters partitioned = + TargetTopicFilter.partition(List.of("eq(attributes/a,1)", "eq(attributes/b,2)")); + + assertThat(partitioned.getRqlExpressions()).containsExactly("eq(attributes/a,1)", "eq(attributes/b,2)"); + assertThat(partitioned.getPipelineExpressions()).isEmpty(); + } + + @Test + public void partitionTrimsEachParam() { + final PartitionedFilters partitioned = TargetTopicFilter.partition( + List.of(" fn:filter(header:ditto-originator,'eq','x')", " gt(attributes/x,5) ")); + + assertThat(partitioned.getPipelineExpressions()) + .containsExactly("fn:filter(header:ditto-originator,'eq','x')"); + assertThat(partitioned.getRqlExpressions()).containsExactly("gt(attributes/x,5)"); + } + + // ===== partition(): RQL params containing "|" or "fn:" stay intact (no splitting anymore) ===== + + @Test + public void partitionKeepsRqlWithPipeInPropertyPathIntact() { + // an unquoted "|" is legal in RQL property paths - without any splitting the whole param stays one + // (valid) RQL expression + assertPureRql("eq(attributes/a|b,1)"); + assertPureRql("exists(attributes/a|b)"); + assertPureRql("like(attributes/x,'*|*')"); + } + + @Test + public void partitionLegacyCombinedSyntaxStaysOneRqlExpression() { + // the retired "|fn:..." single-param syntax is NOT split anymore: not starting with "fn:", the whole + // param is classified as RQL and fails loudly at RQL validation time (ConnectionValidatorTest locks the + // rejection) + final String legacyCombined = "gt(attributes/x,5)|fn:filter(header:a,'exists')"; + final PartitionedFilters partitioned = TargetTopicFilter.partition(List.of(legacyCombined)); + + assertThat(partitioned.getRqlExpressions()).containsExactly(legacyCombined); + assertThat(partitioned.getPipelineExpressions()).isEmpty(); + } + + // ===== partition(): empty / whitespace-only params (a PRESENT empty RQL entry, not absent) ===== + + @Test + public void partitionEmptyParamYieldsPresentEmptyRqlEntryNotAbsent() { + // regression lock: an empty filter param must NOT vanish - it stays a PRESENT (empty) RQL entry, which + // lets ConnectionValidator route it into RQL validation and reject it with InvalidRqlExpressionException, + // exactly as an empty filter was rejected before target topic pipeline filters existed. + assertThat(TargetTopicFilter.partition(List.of("")).getRqlExpressions()).containsExactly(""); + assertThat(TargetTopicFilter.partition(List.of(" ")).getRqlExpressions()).containsExactly(""); + } + + // ===== matchesPipelineFilter(): match / non-match against a signal ===== + + @Test + public void matchesPipelineFilterEqMatchesOnDittoOriginator() { + final Signal signal = thingModifiedWithHeader("ditto-originator", "some:subject"); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'eq','some:subject')", signal, CONNECTION_ID)).isTrue(); + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'eq','other:subject')", signal, CONNECTION_ID)).isFalse(); + } + + @Test + public void matchesPipelineFilterNeMatchesOnDittoOriginator() { + final Signal signal = thingModifiedWithHeader("ditto-originator", "some:subject"); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'ne','other:subject')", signal, CONNECTION_ID)).isTrue(); + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'ne','some:subject')", signal, CONNECTION_ID)).isFalse(); + } + + @Test + public void matchesPipelineFilterOnDittoOriginHeader() { + final Signal signal = thingModifiedWithHeader("ditto-origin", "some-origin"); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-origin,'eq','some-origin')", signal, CONNECTION_ID)).isTrue(); + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-origin,'eq','other-origin')", signal, CONNECTION_ID)).isFalse(); + } + + // ===== matchesPipelineFilter(): absent-header semantics (verified facts, Fact 5) ===== + + @Test + public void matchesPipelineFilterAbsentHeaderEqDrops() { + final Signal signal = thingModifiedWithHeaders(Collections.emptyMap()); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'eq','x')", signal, CONNECTION_ID)).isFalse(); + } + + @Test + public void matchesPipelineFilterAbsentHeaderNePublishes() { + final Signal signal = thingModifiedWithHeaders(Collections.emptyMap()); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'ne','x')", signal, CONNECTION_ID)).isTrue(); + } + + @Test + public void matchesPipelineFilterAbsentHeaderLikeDrops() { + final Signal signal = thingModifiedWithHeaders(Collections.emptyMap()); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'like','some:*')", signal, CONNECTION_ID)).isFalse(); + } + + @Test + public void matchesPipelineFilterAbsentHeaderTwoParamExistsDrops() { + final Signal signal = thingModifiedWithHeaders(Collections.emptyMap()); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'exists')", signal, CONNECTION_ID)).isFalse(); + } + + // ===== matchesPipelineFilter(): chained stages (AND semantics) ===== + + @Test + public void matchesPipelineFilterChainedStagesBothMatchPublishes() { + final Signal signal = thingModifiedWithHeaders(Map.of( + "ditto-originator", "some:subject", + "ditto-origin", "some-origin")); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'eq','some:subject')" + + "|fn:filter(header:ditto-origin,'eq','some-origin')", signal, CONNECTION_ID)).isTrue(); + } + + @Test + public void matchesPipelineFilterChainedStagesFirstNonMatchDrops() { + final Signal signal = thingModifiedWithHeaders(Map.of( + "ditto-originator", "other:subject", + "ditto-origin", "some-origin")); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'eq','some:subject')" + + "|fn:filter(header:ditto-origin,'eq','some-origin')", signal, CONNECTION_ID)).isFalse(); + } + + @Test + public void matchesPipelineFilterChainedStagesSecondNonMatchDrops() { + final Signal signal = thingModifiedWithHeaders(Map.of( + "ditto-originator", "some:subject", + "ditto-origin", "other-origin")); + + assertThat(TargetTopicFilter.matchesPipelineFilter( + "fn:filter(header:ditto-originator,'eq','some:subject')" + + "|fn:filter(header:ditto-origin,'eq','some-origin')", signal, CONNECTION_ID)).isFalse(); + } + + // ===== validatePipelineFilter(): invalid expressions ===== + + @Test + public void validatePipelineFilterAcceptsSingleStage() { + assertThatNoException().isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter(header:ditto-originator,'ne','x')", + DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterRejectsUnknownFunction() { + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class).isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:unknownfn('x')", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterRejectsFilterWithoutArguments() { + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class).isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter()", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterRejectsUnknownPlaceholderPrefix() { + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class).isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter(bogus:x,'eq','y')", DittoHeaders.empty())); + } + + // ===== validatePipelineFilter(): chained stages, pipeline grammar limits ===== + + @Test + public void validatePipelineFilterAcceptsChainedStages() { + assertThatNoException().isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter( + "fn:filter(header:a,'exists')|fn:filter(header:b,'exists')", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterRejectsChainedParamWithUnknownFunctionStage() { + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class).isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter( + "fn:filter(header:a,'exists')|fn:unknownfn('x')", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterAcceptsTenChainedStagesButRejectsEleven() { + // the resolver's pipeline grammar caps a pipeline at 10 fn: stages (the internal fn:default seed does + // not eat into the user's budget: seed + 10 user stages is exactly the grammar's 11-element maximum) + final String tenStages = String.join("|", Collections.nCopies(10, "fn:filter(header:a,'exists')")); + assertThatNoException().isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter(tenStages, DittoHeaders.empty())); + + final String elevenStages = String.join("|", Collections.nCopies(11, "fn:filter(header:a,'exists')")); + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class).isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter(elevenStages, DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterAcceptsQuotedPipeInsideChainedStages() { + // the resolver's stage split is quote-aware: the '|' inside 'a|b' must not be taken for a stage separator + assertThatNoException().isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter( + "fn:filter(header:a,'eq','a|b')|fn:filter(header:b,'exists')", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterRejectsTrailingPipe() { + // rejected by the resolver's pipeline grammar (empty trailing stage), no custom scan involved + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class).isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter(header:a,'exists')|", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterAcceptsQuotedPipeInSingleQuotedConstant() { + assertThatNoException().isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter(header:a,'eq','a|b')", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterAcceptsQuotedPipeInDoubleQuotedConstant() { + assertThatNoException().isThrownBy(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter(header:a,'eq',\"a|b\")", DittoHeaders.empty())); + } + + @Test + public void validatePipelineFilterTrailingBackslashDoesNotThrowUnexpectedly() { + // a trailing backslash must never escape the documented exception contract; the resolver validation may + // still reject the expression, but only ever with the documented exception type + final Throwable throwable = catchThrowable(() -> + TargetTopicFilter.validatePipelineFilter("fn:filter(header:a,'exists')\\", DittoHeaders.empty())); + + if (throwable != null) { + assertThat(throwable).isInstanceOf(ConnectionConfigurationInvalidException.class); + } + } + + // ===== test helpers ===== + + private static Signal thingModifiedWithHeader(final String key, final String value) { + return thingModifiedWithHeaders(Map.of(key, value)); + } + + private static Signal thingModifiedWithHeaders(final Map headers) { + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) + .build(); + final DittoHeaders dittoHeaders = DittoHeaders.newBuilder().putHeaders(headers).build(); + return ThingModified.of(thing, 1L, Instant.now(), dittoHeaders, null); + } +} diff --git a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilterWithFilterTest.java b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilterWithFilterTest.java index 9cf1b80fa67..816492ba6b2 100644 --- a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilterWithFilterTest.java +++ b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/persistence/SignalFilterWithFilterTest.java @@ -37,6 +37,7 @@ import org.eclipse.ditto.connectivity.model.HeaderMapping; import org.eclipse.ditto.connectivity.model.Target; import org.eclipse.ditto.connectivity.service.messaging.TestConstants; +import org.eclipse.ditto.connectivity.service.messaging.monitoring.ConnectionMonitor; import org.eclipse.ditto.connectivity.service.messaging.monitoring.ConnectionMonitorRegistry; import org.eclipse.ditto.json.JsonPointer; import org.eclipse.ditto.json.JsonValue; @@ -50,6 +51,7 @@ import org.eclipse.ditto.things.model.ThingId; import org.eclipse.ditto.things.model.signals.events.ThingModified; import org.junit.Test; +import org.mockito.Mockito; /** * Tests {@link SignalFilter} for filtering with namespace + RQL filter. @@ -345,4 +347,543 @@ public void applySignalFilterWithFeatureIdPlaceholder() { final List filteredTargets = signalFilter.filter(signal); Assertions.assertThat(filteredTargets).hasSize(1).contains(target); } + + // ===== pure pipeline (fn:) target topic filter ===== + + @Test + public void applySignalFilterWithPurePipelineFilterMatchesAndNonMatchesOnDittoOriginator() { + final String filter = "fn:filter(header:ditto-originator,'eq','some:subject')"; + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilter(filter) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) + .build(); + + final DittoHeaders matchingHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "some:subject") + .build(); + final ThingModified matching = ThingModified.of(thing, 3L, Instant.now(), matchingHeaders, null); + + final DittoHeaders nonMatchingHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .build(); + final ThingModified nonMatching = ThingModified.of(thing, 3L, Instant.now(), nonMatchingHeaders, null); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + assertThat(signalFilter.filter(matching)).containsOnly(target); // THEN: matching originator is returned + assertThat(signalFilter.filter(nonMatching)).isEmpty(); // THEN: non-matching originator is not returned + } + + @Test + public void applySignalFilterWithPurePipelineFilterAbsentHeaderNegationPublishes() { + // absent header + "ne" => publish (verified fact 5) + final String filter = "fn:filter(header:ditto-originator,'ne','some:subject')"; + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilter(filter) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) + .build(); + final DittoHeaders headers = DittoHeaders.newBuilder() // WHEN: no "ditto-originator" header is set + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .build(); + final ThingModified thingModified = ThingModified.of(thing, 3L, Instant.now(), headers, null); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + assertThat(signalFilter.filter(thingModified)).containsOnly(target); + } + + // ===== RQL and pipeline filter params on one topic (AND semantics) ===== + + @Test + public void applySignalFilterWithRqlAndPipelineFilterParams() { + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilters(List.of("eq(attributes/test,42)", + "fn:filter(header:ditto-originator,'eq','some:subject')")) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing matchingThing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) // RQL matches + .build(); + final Thing nonMatchingThing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(99)) // RQL does not match + .build(); + + final DittoHeaders matchingOriginatorHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "some:subject") // pipeline matches + .build(); + final DittoHeaders nonMatchingOriginatorHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") // pipeline does not match + .build(); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + // RQL matches + pipeline matches => returned + final ThingModified rqlMatchPipelineMatch = + ThingModified.of(matchingThing, 3L, Instant.now(), matchingOriginatorHeaders, null); + assertThat(signalFilter.filter(rqlMatchPipelineMatch)).containsOnly(target); + + // RQL matches + pipeline does not match => not returned + final ThingModified rqlMatchPipelineNonMatch = + ThingModified.of(matchingThing, 3L, Instant.now(), nonMatchingOriginatorHeaders, null); + assertThat(signalFilter.filter(rqlMatchPipelineNonMatch)).isEmpty(); + + // RQL does not match + pipeline matches => not returned + final ThingModified rqlNonMatchPipelineMatch = + ThingModified.of(nonMatchingThing, 3L, Instant.now(), matchingOriginatorHeaders, null); + assertThat(signalFilter.filter(rqlNonMatchPipelineMatch)).isEmpty(); + } + + // ===== T-I2: pure pipeline filter for LIVE_MESSAGES (primary echo-suppression use case) ===== + + @Test + public void applySignalFilterForLiveMessagesWithPurePipelineFilter() { + final String filter = "fn:filter(header:ditto-originator,'eq','some:subject')"; + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("message/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(LIVE_MESSAGES) + .withFilter(filter) + .build()) + .build(); + + final Connection connection = + ConnectivityModelFactory.newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, + ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final DittoHeaders matchingHeaders = DittoHeaders.newBuilder() + .channel(TopicPath.Channel.LIVE.getName()) + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "some:subject") + .build(); + final SendThingMessage matchingMessage = SendThingMessage.of(THING_ID, + Message.newBuilder(MessageHeaders.newBuilder(MessageDirection.TO, THING_ID, "fubar").build()) + .build(), matchingHeaders); + + final DittoHeaders nonMatchingHeaders = DittoHeaders.newBuilder() + .channel(TopicPath.Channel.LIVE.getName()) + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .build(); + final SendThingMessage nonMatchingMessage = SendThingMessage.of(THING_ID, + Message.newBuilder(MessageHeaders.newBuilder(MessageDirection.TO, THING_ID, "fubar").build()) + .build(), nonMatchingHeaders); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + assertThat(signalFilter.filter(matchingMessage)).containsOnly(target); + assertThat(signalFilter.filter(nonMatchingMessage)).isEmpty(); + } + + // ===== T-I3: OR across two FilteredTopics (same base Topic, different filter) on ONE target ===== + + @Test + public void applySignalFilterOrAcrossTwoFilteredTopicsOnOneTarget() { + // topic A: non-matching RQL filter + final String rqlFilter = "eq(attributes/test,999)"; + // topic B: matching pipeline filter + final String pipelineFilter = "fn:filter(header:ditto-originator,'eq','some:subject')"; + + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS).withFilter(rqlFilter).build(), + ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS).withFilter(pipelineFilter) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) // never matches topic A's RQL filter + .build(); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + // topic A (RQL) does not match, topic B (pipeline) matches => returned + final DittoHeaders matchingHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "some:subject") + .build(); + final ThingModified matching = ThingModified.of(thing, 3L, Instant.now(), matchingHeaders, null); + assertThat(signalFilter.filter(matching)).containsOnly(target); + + // neither topic A (RQL) nor topic B (pipeline) matches => not returned + final DittoHeaders nonMatchingHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .build(); + final ThingModified nonMatching = ThingModified.of(thing, 3L, Instant.now(), nonMatchingHeaders, null); + assertThat(signalFilter.filter(nonMatching)).isEmpty(); + } + + // ===== T-I4: mix of a pure-pipeline target and a pure-RQL target on one connection ===== + + @Test + public void applySignalFilterWithMixOfPipelineTargetAndRqlTarget() { + final String pipelineFilter = "fn:filter(header:ditto-originator,'eq','some:subject')"; + final Target targetA = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilter(pipelineFilter) + .build()) + .build(); + + final String rqlFilter = "eq(attributes/test,42)"; + final Target targetB = ConnectivityModelFactory.newTargetBuilder() + .address("twin/b") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilter(rqlFilter) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(targetA, targetB)) + .build(); + + final Thing matchingThing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) // matches targetB's RQL filter + .build(); + final Thing nonMatchingThing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(99)) // does not match targetB's RQL filter + .build(); + + final DittoHeaders matchingOriginatorHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "some:subject") // matches targetA's pipeline filter + .build(); + final DittoHeaders nonMatchingOriginatorHeaders = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") // does not match targetA's pipeline filter + .build(); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + // both match + assertThat(signalFilter.filter( + ThingModified.of(matchingThing, 3L, Instant.now(), matchingOriginatorHeaders, null))) + .containsOnly(targetA, targetB); + + // only targetB (RQL) matches + assertThat(signalFilter.filter( + ThingModified.of(matchingThing, 3L, Instant.now(), nonMatchingOriginatorHeaders, null))) + .containsOnly(targetB); + + // only targetA (pipeline) matches + assertThat(signalFilter.filter( + ThingModified.of(nonMatchingThing, 3L, Instant.now(), matchingOriginatorHeaders, null))) + .containsOnly(targetA); + + // neither matches + assertThat(signalFilter.filter( + ThingModified.of(nonMatchingThing, 3L, Instant.now(), nonMatchingOriginatorHeaders, null))) + .isEmpty(); + } + + // ===== T-M5: regression - pure RQL filter with a literal "|" in a quoted attribute value ===== + + @Test + public void applySignalFilterWithRqlFilterContainingLiteralPipeInAttributeValue() { + final String filter = "like(attributes/test,'*|*')"; + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilter(filter) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of("a|b")) // WHEN: literal "|" in the value + .build(); + final DittoHeaders headers = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .build(); + final ThingModified thingModified = ThingModified.of(thing, 3L, Instant.now(), headers, null); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + assertThat(signalFilter.filter(thingModified)).containsOnly(target); + } + + // ===== runtime failure policy: pipeline evaluation failure drops the target AND records a connection-log entry ===== + + @Test + public void applySignalFilterWithFailingPipelineFilterDropsTargetAndRecordsConnectionLogFailure() { + // parses fine as a pure pipeline filter but throws a DittoRuntimeException at evaluation time + // (unknown pipeline function) + final String filter = "fn:unknownfn('x')"; + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilter(filter) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) + .build(); + final DittoHeaders headers = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .build(); + final ThingModified thingModified = ThingModified.of(thing, 3L, Instant.now(), headers, null); + + // a local (non-stub-only) registry so that the failure interaction can be verified + final ConnectionMonitor filteredMonitor = Mockito.mock(ConnectionMonitor.class); + @SuppressWarnings("unchecked") + final ConnectionMonitorRegistry registry = Mockito.mock(ConnectionMonitorRegistry.class); + Mockito.when(registry.forOutboundDispatched(Mockito.any(Connection.class), Mockito.anyString())) + .thenReturn(Mockito.mock(ConnectionMonitor.class)); + Mockito.when(registry.forOutboundFiltered(Mockito.any(Connection.class), Mockito.anyString())) + .thenReturn(filteredMonitor); + + final SignalFilter signalFilter = new SignalFilter(connection, registry); + + // THEN: the evaluation failure is treated as a non-match ... + assertThat(signalFilter.filter(thingModified)).isEmpty(); + // ... AND recorded as a FAILURE entry in the target's user-visible connection log + Mockito.verify(filteredMonitor).failure(Mockito.eq(thingModified), Mockito.anyString(), + Mockito.eq(filter), Mockito.anyString()); + } + + // ===== chained pipeline stages in one filter param (AND semantics) ===== + + @Test + public void applySignalFilterWithChainedPipelineStagesAndSemantics() { + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilters(List.of( + "fn:filter(header:ditto-originator,'ne','excluded:subject')" + + "|fn:filter(header:ditto-origin,'ne','excluded-connection')")) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) + .build(); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + // both chained stages match => returned + final DittoHeaders bothMatch = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .putHeader("ditto-origin", "other-connection") + .build(); + assertThat(signalFilter.filter(ThingModified.of(thing, 3L, Instant.now(), bothMatch, null))) + .containsOnly(target); + + // first chained stage does not match => not returned + final DittoHeaders firstNonMatch = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "excluded:subject") + .putHeader("ditto-origin", "other-connection") + .build(); + assertThat(signalFilter.filter(ThingModified.of(thing, 3L, Instant.now(), firstNonMatch, null))) + .isEmpty(); + + // second chained stage does not match => not returned + final DittoHeaders secondNonMatch = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .putHeader("ditto-origin", "excluded-connection") + .build(); + assertThat(signalFilter.filter(ThingModified.of(thing, 3L, Instant.now(), secondNonMatch, null))) + .isEmpty(); + } + + @Test + public void applySignalFilterWithRqlAndChainedPipelineFilterParam() { + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilters(List.of( + "eq(attributes/test,42)", + "fn:filter(header:ditto-originator,'ne','excluded:subject')" + + "|fn:filter(header:ditto-origin,'ne','excluded-connection')")) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing matchingThing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) // RQL matches + .build(); + final Thing nonMatchingThing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(99)) // RQL does not match + .build(); + final DittoHeaders allPipelinesMatch = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .putHeader("ditto-origin", "other-connection") + .build(); + final DittoHeaders firstPipelineNonMatch = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "excluded:subject") + .putHeader("ditto-origin", "other-connection") + .build(); + final DittoHeaders secondPipelineNonMatch = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .putHeader("ditto-origin", "excluded-connection") + .build(); + + final SignalFilter signalFilter = new SignalFilter(connection, connectionMonitorRegistry); + + // RQL + both chained stages match => returned + assertThat(signalFilter.filter( + ThingModified.of(matchingThing, 3L, Instant.now(), allPipelinesMatch, null))) + .containsOnly(target); + + // RQL does not match => not returned + assertThat(signalFilter.filter( + ThingModified.of(nonMatchingThing, 3L, Instant.now(), allPipelinesMatch, null))) + .isEmpty(); + + // first chained stage does not match => not returned + assertThat(signalFilter.filter( + ThingModified.of(matchingThing, 3L, Instant.now(), firstPipelineNonMatch, null))) + .isEmpty(); + + // second chained stage does not match => not returned + assertThat(signalFilter.filter( + ThingModified.of(matchingThing, 3L, Instant.now(), secondPipelineNonMatch, null))) + .isEmpty(); + } + + @Test + public void applySignalFilterWithFailingSecondPipelineFilterParamRecordsFailureForThatParam() { + // DEFENSIVE lock: validation now caps a topic at one pipeline filter param, but the runtime AND-loop + // must keep handling multi-param lists that never passed validation (model-built or pre-rule persisted + // connections). The first param evaluates fine (and matches), the second throws at evaluation time - + // the failure entry must name the FAILING param, not the first one. + final String matchingFilter = "fn:filter(header:ditto-originator,'ne','excluded:subject')"; + final String failingFilter = "fn:unknownfn('x')"; + final Target target = ConnectivityModelFactory.newTargetBuilder() + .address("twin/a") + .authorizationContext(newAuthContext(DittoAuthorizationContextType.UNSPECIFIED, AUTHORIZED)) + .headerMapping(HEADER_MAPPING) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(TWIN_EVENTS) + .withFilters(List.of(matchingFilter, failingFilter)) + .build()) + .build(); + + final Connection connection = ConnectivityModelFactory + .newConnectionBuilder(CONNECTION_ID, ConnectionType.AMQP_10, ConnectivityStatus.OPEN, URI) + .targets(List.of(target)) + .build(); + + final Thing thing = Thing.newBuilder() + .setId(THING_ID) + .setAttribute(JsonPointer.of("test"), JsonValue.of(42)) + .build(); + final DittoHeaders headers = DittoHeaders.newBuilder() + .readGrantedSubjects(Collections.singletonList(AUTHORIZED)) + .putHeader("ditto-originator", "other:subject") + .build(); + final ThingModified thingModified = ThingModified.of(thing, 3L, Instant.now(), headers, null); + + final ConnectionMonitor filteredMonitor = Mockito.mock(ConnectionMonitor.class); + @SuppressWarnings("unchecked") + final ConnectionMonitorRegistry registry = Mockito.mock(ConnectionMonitorRegistry.class); + Mockito.when(registry.forOutboundDispatched(Mockito.any(Connection.class), Mockito.anyString())) + .thenReturn(Mockito.mock(ConnectionMonitor.class)); + Mockito.when(registry.forOutboundFiltered(Mockito.any(Connection.class), Mockito.anyString())) + .thenReturn(filteredMonitor); + + final SignalFilter signalFilter = new SignalFilter(connection, registry); + + assertThat(signalFilter.filter(thingModified)).isEmpty(); + Mockito.verify(filteredMonitor).failure(Mockito.eq(thingModified), Mockito.anyString(), + Mockito.eq(failingFilter), Mockito.anyString()); + } } diff --git a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidatorTest.java b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidatorTest.java index 8bbfb907229..3397e9f24b2 100644 --- a/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidatorTest.java +++ b/connectivity/service/src/test/java/org/eclipse/ditto/connectivity/service/messaging/validation/ConnectionValidatorTest.java @@ -36,6 +36,7 @@ import org.eclipse.ditto.base.model.acks.AcknowledgementLabel; import org.eclipse.ditto.base.model.acks.AcknowledgementLabelInvalidException; import org.eclipse.ditto.base.model.acks.AcknowledgementLabelNotUniqueException; +import org.eclipse.ditto.base.model.exceptions.InvalidRqlExpressionException; import org.eclipse.ditto.base.model.headers.DittoHeaders; import org.eclipse.ditto.base.model.json.Jsonifiable; import org.eclipse.ditto.connectivity.model.ClientCertificateCredentials; @@ -453,6 +454,234 @@ public void acceptValidConnectionWithValidTargetFilterContainingPlaceholders() { underTest.validate(connection, DittoHeaders.empty(), actorSystem); } + @Test + public void acceptValidConnectionWithPurePipelineTargetFilter() { + final List targetWithValidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:filter(header:ditto-originator,'eq','some:subject')") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithValidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + underTest.validate(connection, DittoHeaders.empty(), actorSystem); + } + + @Test + public void acceptValidConnectionWithRqlAndPipelineTargetFilterParams() { + final List targetWithValidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of("eq(attributes/a,1)", + "fn:filter(header:ditto-originator,'eq','some:subject')")) + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithValidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + underTest.validate(connection, DittoHeaders.empty(), actorSystem); + } + + @Test + public void rejectConnectionWithTwoPipelineTargetFilterParams() { + // at most one of a topic's filter params may be a pipeline expression - several pipeline conditions + // belong into ONE fn: param, chained with '|' + final List targetWithInvalidFilters = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of( + "fn:filter(header:ditto-originator,'ne','some:subject')", + "fn:filter(header:ditto-origin,'ne','some-connection-id')")) + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithInvalidFilters) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)) + .withMessageContaining("at most one pipeline filter"); + } + + @Test + public void rejectConnectionWithTwoRqlTargetFilterParams() { + // at most one of a topic's filter params may be an RQL expression - several RQL conditions belong into a + // single RQL expression combined with and(...) + final List targetWithInvalidFilters = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of("eq(attributes/a,1)", "eq(attributes/b,2)")) + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithInvalidFilters) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)) + .withMessageContaining("at most one RQL filter"); + } + + @Test + public void acceptConnectionWithChainedPipelineTargetFilterParam() { + // several fn: stages chained with '|' inside the ONE pipeline filter param are the intended way to + // AND several pipeline conditions + final List targetWithValidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:filter(header:a,'exists')|fn:filter(header:b,'exists')") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithValidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + underTest.validate(connection, DittoHeaders.empty(), actorSystem); + } + + @Test + public void rejectConnectionWithLegacyCombinedFilterSyntax() { + // locks the loud failure of the retired "|fn:..." single-param syntax: not starting with "fn:", the + // whole param is routed into the RQL parser, which rejects the "|fn:..." tail - it must NOT be silently + // split or accepted. + final List targetWithLegacyFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter( + "gt(attributes/counter,42)|fn:filter(header:ditto-originator,'eq','x')") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithLegacyFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(InvalidRqlExpressionException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)); + } + + @Test + public void acceptValidConnectionWithLegacyTargetFilterContainingUnquotedPipeInPropertyPath() { + // back-compat lock: "eq(attributes/a|b,1)" is a pre-existing, valid pure RQL expression (an unquoted "|" + // is legal in RQL property paths) and must keep being accepted verbatim - filter params are classified + // only by their "fn:" prefix, never split at a "|". + final List targetWithValidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("eq(attributes/a|b,1)") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithValidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + underTest.validate(connection, DittoHeaders.empty(), actorSystem); + } + + @Test + public void rejectConnectionWithTargetFilterContainingUnknownPipelineFunction() { + final List targetWithInvalidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("fn:unknownfn('x')") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithInvalidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(ConnectionConfigurationInvalidException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)); + } + + @Test + public void rejectConnectionWithMalformedPureRqlTargetFilterAsInvalidRqlExpression() { + // locks C-I1: a malformed pure RQL filter (no "fn:" involved at all) must keep failing with + // InvalidRqlExpressionException straight from the RQL parser, and must NOT be misrouted into + // ConnectionConfigurationInvalidException (the pipeline-validation exception type). + final List targetWithInvalidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("gt(attributes/x,)") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithInvalidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(InvalidRqlExpressionException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)); + } + + @Test + public void rejectConnectionWithMalformedRqlFilterParamAlongsideValidPipelineFilterParam() { + // a malformed RQL filter param keeps failing with InvalidRqlExpressionException even when the topic also + // carries a (valid) pipeline filter param + final List targetWithInvalidFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilters(List.of("gt(attributes/x,)", + "fn:filter(header:ditto-originator,'eq','x')")) + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithInvalidFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(InvalidRqlExpressionException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)); + } + + @Test + public void rejectConnectionWithEmptyTargetFilterAsInvalidRqlExpression() { + // regression lock: an empty target topic filter string (e.g. a target address ending in "?filter=") must + // keep being rejected with InvalidRqlExpressionException, exactly as before target topic pipeline filters + // existed - TargetTopicFilter.partition classifies an empty param as a PRESENT (empty) RQL entry, which + // is routed into RQL validation here rather than silently accepted. + final List targetWithEmptyFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter("") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithEmptyFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(InvalidRqlExpressionException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)); + } + + @Test + public void rejectConnectionWithWhitespaceOnlyTargetFilterAsInvalidRqlExpression() { + final List targetWithBlankFilter = singletonList( + ConnectivityModelFactory.newTargetBuilder(TestConstants.Targets.TWIN_TARGET) + .topics(ConnectivityModelFactory.newFilteredTopicBuilder(Topic.TWIN_EVENTS) + .withFilter(" ") + .build()) + .build()); + final Connection connection = createConnection(CONNECTION_ID) + .toBuilder() + .setTargets(targetWithBlankFilter) + .build(); + final ConnectionValidator underTest = getConnectionValidator(); + assertThatExceptionOfType(InvalidRqlExpressionException.class) + .isThrownBy(() -> underTest.validate(connection, DittoHeaders.empty(), actorSystem)); + } + @Test public void acceptValidConnectionWithValidNumberPayloadMapping() { final Connection connection = createConnection(CONNECTION_ID) diff --git a/documentation/src/main/resources/jsonschema/connection.json b/documentation/src/main/resources/jsonschema/connection.json index e9fe2d3ece8..63d8a31798f 100644 --- a/documentation/src/main/resources/jsonschema/connection.json +++ b/documentation/src/main/resources/jsonschema/connection.json @@ -656,11 +656,13 @@ "items": { "type": "string", "title": "Subscribed topics.", - "description": "Contains the type of messages that are delivered to this target. You can receive\n * Thing events: `_/_/things/twin/events` (notification about twin change) \n * Live events: `_/_/things/live/events`\n * Live commands: `_/_/things/live/commands`\n * Live messages: `_/_/things/live/messages`\n\nYou can specify an additional namespace and/or event filter (URL encoded)", + "description": "Contains the type of messages that are delivered to this target. You can receive\n * Thing events: `_/_/things/twin/events` (notification about twin change) \n * Live events: `_/_/things/live/events`\n * Live commands: `_/_/things/live/commands`\n * Live messages: `_/_/things/live/messages`\n\nYou can specify an additional namespace and/or event filter (URL encoded). The `filter` parameter may be repeated: at most one RQL expression and at most one placeholder pipeline expression starting with `fn:` (several `fn:` stages chainable with `|`); all given filters must match (AND)", "examples": [ "_/_/things/twin/events", "_/_/things/twin/events?namespaces=org.eclipse.ditto.one,org.eclipse.foo", "_/_/things/twin/events?namespaces=org.eclipse.ditto&filter=eq(attributes/counter,42)", + "_/_/things/twin/events?filter=fn:filter(header:ditto-originator,'ne','some:excluded-subject')", + "_/_/things/twin/events?filter=gt(attributes/counter,42)&filter=fn:filter(header:ditto-originator,'ne','some:excluded-subject')", "_/_/things/twin/events?extraFields=attributes", "_/_/things/twin/events?extraFields=attributes&filter=eq(attributes/counter,42)", "_/_/things/live/commands", diff --git a/documentation/src/main/resources/pages/ditto/basic-changenotifications.md b/documentation/src/main/resources/pages/ditto/basic-changenotifications.md index c5a7dbac22a..45c66c3c5f8 100644 --- a/documentation/src/main/resources/pages/ditto/basic-changenotifications.md +++ b/documentation/src/main/resources/pages/ditto/basic-changenotifications.md @@ -52,6 +52,11 @@ For more granular control, use an [RQL expression](basic-rql.html) to filter bas {% include note.html content="The RQL filter applies to the *modified* data by default. Unchanged data is only considered when it has been [enriched via extraFields](basic-enrichment.html)." %} +[Connections](basic-connections.html) additionally accept a placeholder function pipeline (`fn:...`) +in an additional `filter` parameter alongside (or instead of) an RQL expression -- see +[Filtering with placeholder functions](basic-connections.html#filtering-with-placeholder-functions). +This is not available for the WebSocket API or SSE, which only support a single RQL filter. + ### Examples Only emit events when `count` changes to a value greater than 42: diff --git a/documentation/src/main/resources/pages/ditto/basic-connections.md b/documentation/src/main/resources/pages/ditto/basic-connections.md index 493d1b73ab6..b7429dc36b3 100644 --- a/documentation/src/main/resources/pages/ditto/basic-connections.md +++ b/documentation/src/main/resources/pages/ditto/basic-connections.md @@ -292,8 +292,14 @@ You define which message types to publish via the `topics` array. You can filter | `_/_/policies/announcements` | ✔ | ❌ | | `_/_/connections/announcements` | ❌ | ❌ | -Filter parameters use HTTP query parameter syntax (`?` for the first, `&` for subsequent). URL-encode -filter values before using them: +Filter parameters use HTTP query parameter syntax (`?` for the first, `&` for subsequent). The +`filter` parameter may be given **twice** on one topic -- at most one RQL expression and at most one +placeholder pipeline expression (`fn:...`); all given filters must match for a signal to be +published (**AND** semantics, see +[filtering with placeholder functions](#filtering-with-placeholder-functions) below). URL-encode +filter values before using them: topic filters given in this string form are URL-decoded when parsed, +so a literal `+` (decoded to a space) or `%xx` sequence in a compared value must itself be +URL-encoded -- this applies to RQL `like` patterns and pipeline compared values alike: ```json { @@ -307,6 +313,123 @@ filter values before using them: } ``` +If a target's `topics` array lists several topic entries, they are evaluated independently and +combined with **OR** semantics -- a signal is published as soon as it matches *any one* listed +topic (each with its own namespace/RQL/pipeline filter). + +### Filtering with placeholder functions + +In addition to (or instead of) an RQL expression, a target topic may carry `filter` parameters +holding a placeholder function invocation from the +[function library](basic-placeholders.html#function-library), most commonly +[`fn:filter()`](basic-placeholders.html#function-library). Such a pipeline expression is +evaluated per outbound signal against that signal's headers, topic, entity, and time -- see +[connection target topic filter placeholders](basic-placeholders.html#scope-connection-target-topic-filter) +for the full list -- instead of against thing/event *data*. Because of that, a pipeline filter also works for topics for which an +RQL filter cannot meaningfully match, such as `_/_/things/live/commands` (marked ❌ for +"RQL filter" in the table above). + +A pipeline expression: +* **resolves** (produces a value) -- the target topic is **published** +* stays **unresolved** (the filter drops the value) -- the target topic is **suppressed** + +The primary use case is suppressing events caused by a given subject, or caused by another +connection. Each is a standalone filter (do **not** combine them as two separate `topics` entries -- +that would be an OR, publishing whenever *either* condition holds; see below on how to combine +conditions with AND): + +```json +{ + "address": "", + "topics": [ + "_/_/things/twin/events?filter=fn:filter(header:ditto-originator,'ne','some:excluded-subject')" + ], + "authorizationContext": ["ditto:outbound-auth-subject"] +} +``` + +* `header:ditto-originator` resolves to the first authorization subject of the request that caused + the signal. +* `header:ditto-origin` resolves to the ID of the connection that originally caused the signal, e.g. + `filter=fn:filter(header:ditto-origin,'ne','some-other-connection-id')`. + Ditto already suppresses signals a connection caused itself by default; filtering on + `ditto-origin` is only needed to additionally exclude signals caused by *other* connections. + +A `filter` query parameter whose (trimmed) value starts with `fn:` is a pipeline filter; anything +else is treated as RQL. Several `fn:` stages can be chained with `|` inside the pipeline filter. +Each stage only runs if the previous one resolved (matched), so chaining is **AND**: every stage +must match for the pipeline to resolve, for example: + +```text +filter=fn:filter(header:ditto-originator,'ne','some:subject')|fn:filter(header:ditto-origin,'ne','some-connection-id') +``` + +An RQL filter and a pipeline filter can be combined as two `filter` parameters, again with **AND** +semantics -- the RQL expression and the pipeline must both match: + +```text +filter=gt(attributes/counter,42)&filter=fn:filter(header:ditto-originator,'ne','some:subject') +``` + +At most **one** of a topic's `filter` parameters may be an RQL expression -- combine several RQL +conditions into a single expression with `and(...)` instead. Likewise at most **one** may be a +pipeline expression -- combine several pipeline conditions by chaining `fn:` stages with `|` inside +it. + +#### Absent header behavior + +Since the pipeline evaluates per signal, a referenced header may be absent for a given signal (for +example, `ditto-originator` is absent for signals with no authenticated causing subject, such as +those from Ditto's own internal processing paths). The outcome then depends on the `rqlFunction`: + +| `rqlFunction` | Outcome when the header is absent | +|---------------|------------------------------------| +| `eq` | dropped (suppressed) | +| `ne` | **published** | +| `like` | dropped, unless the pattern itself matches the empty string (e.g. `'*'`) | +| `exists` (2-param form, e.g. `fn:filter(header:x,'exists')`) | dropped | + +{% include important.html content="`ne` on an absent header resolves to **published**, not +suppressed -- this is the opposite of what `eq` does and easy to get wrong. For example, +`fn:filter(header:ditto-originator,'ne','some:subject')` also publishes any signal that never +carries a `ditto-originator` header at all -- because 'absent' trivially satisfies 'not equal to +some:subject'. If only signals that actually carry the header should be affected, combine the +pipeline with an additional `exists` stage." additionalStyle="" %} + +#### Restrictions + +* A pipeline filter parameter must always start with `fn:` -- a bare leading placeholder + (e.g. `filter=header:ditto-originator`) is **not** a valid pipeline filter; place the placeholder + as a function parameter instead, as shown above. +* As with any placeholder pipeline, every stage after the first must itself be an `fn:` function + call -- a bare placeholder cannot appear mid-pipeline. +* A pipeline filter may contain at most **10** `fn:` stages; exceeding the limit is rejected at + connection creation/update time. +* A topic accepts at most **one** RQL `filter` parameter and at most **one** pipeline `filter` + parameter; a second one of either kind is rejected at connection creation/update time. Combine + RQL conditions with `and(...)` and pipeline conditions by chaining stages with `|`. +* An unrecognized `rqlFunction` name (i.e. anything other than `eq`, `ne`, `like`, `exists`) is + **not** rejected at connection creation/update time -- that filter simply never matches at + runtime. Double-check spelling. +* Pipeline placeholders never see fields added via [`extraFields` + enrichment](#target-topics-and-enrichment) -- they only ever see the signal's own headers, topic, + entity, and time. Unlike RQL, a pipeline filter cannot filter on enriched, unchanged data. +* Pipeline filters should only reference headers that are stable for the signal's lifetime (such as + `ditto-originator` or `ditto-origin`): internal bookkeeping headers such as `requested-acks` are + mutated while the signal is processed and are not reliable filter inputs. +* As with RQL, no `filter` at all -- pipeline, RQL, or combined -- can be set on + `_/_/policies/announcements` or `_/_/connections/announcements` (see the table above); it is + silently ignored if present. +* A filter-suppressed signal produces no user-visible log entry, the same as with a pure RQL filter + today. Only when evaluating a pipeline filter *fails* (rather than simply not matching) is a + failure entry recorded in the [connection logs](connectivity-manage-connections.html#connection-logs). + +{% include warning.html content="Only start using `fn:...` filters -- and in particular **repeated** +`filter` parameters -- once **all** instances of your connectivity service run a Ditto version that +supports them. Older instances cannot even *parse* a topic string carrying more than one `filter` +parameter: a connection persisted with the new syntax breaks connection loading on such instances, +it is not merely rejected." %} + ### Target topics and enrichment You can add extra fields to outgoing messages with the `extraFields` parameter. diff --git a/documentation/src/main/resources/pages/ditto/basic-placeholders.md b/documentation/src/main/resources/pages/ditto/basic-placeholders.md index 58dbaa51dd6..8f067cdc15a 100644 --- a/documentation/src/main/resources/pages/ditto/basic-placeholders.md +++ b/documentation/src/main/resources/pages/ditto/basic-placeholders.md @@ -325,6 +325,31 @@ _org.eclipse.ditto/device-123/things/live/messages/hello.world_ these placeholde | `topic:subject` | _hello.world_ | | `topic:action-subject` | _hello.world_ | +### Scope: Connection target topic filter + +In a connection's [target topic filter](basic-connections.html#filtering-with-placeholder-functions), +a placeholder function pipeline (most commonly [`fn:filter()`](#function-library), several stages +chainable with `|`) may be used as a `filter` query parameter, alongside an optional +[RQL expression](basic-rql.html) `filter` parameter (at most one of each; all of a topic's filters +are combined with AND). As with +[RQL expressions when filtering for Ditto Protocol messages](#scope-rql-expressions-when-filtering-for-ditto-protocol-messages), +such a pipeline is a bare expression and placeholders must not be surrounded by curly braces, e.g. +`fn:filter(header:ditto-originator,'ne','some:subject')`. The pipeline is evaluated per outbound +signal, before enrichment, so the following placeholders are available in general: +* [entity placeholder](#entity-placeholder) +* [thing placeholder](#thing-placeholder) +* [thing-json placeholder](#thing-json-placeholder) +* [feature placeholder](#feature-placeholder) +* [header placeholder](#header-placeholder) +* [request placeholder](#request-placeholder) +* [resource placeholder](#resource-placeholder) +* [topic placeholder](#topic-placeholder) +* [time placeholder](#time-placeholder) +* [connection placeholder](#connection-placeholder) + +Unlike the [Connections](#scope-connections) scope used e.g. for target addresses, these +placeholders never see fields declared via `extraFields` enrichment, and no +[policy placeholder](#policy-placeholder) is available. ## Function expressions