This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 03e3ef9cb71b CAMEL-25253: camel-dynamic-router - an update of a
subscription must replace it (#27250)
03e3ef9cb71b is described below
commit 03e3ef9cb71bcf30b25aced39a4292646fce50e9
Author: allthingssecurity <[email protected]>
AuthorDate: Fri Oct 2 13:23:04 2026 +0530
CAMEL-25253: camel-dynamic-router - an update of a subscription must
replace it (#27250)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../filter/DynamicRouterFilterService.java | 30 +++++--
.../DynamicRouterFilterServiceUpdateTest.java | 91 ++++++++++++++++++++++
2 files changed, 113 insertions(+), 8 deletions(-)
diff --git
a/components/camel-dynamic-router/src/main/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterService.java
b/components/camel-dynamic-router/src/main/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterService.java
index 8664959d3a52..dd4c282b0962 100644
---
a/components/camel-dynamic-router/src/main/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterService.java
+++
b/components/camel-dynamic-router/src/main/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterService.java
@@ -170,17 +170,31 @@ public class DynamicRouterFilterService {
* @return the ID of the added filter
*/
public String addFilterForChannel(final PrioritizedFilter filter, final
String channel, final boolean update) {
- boolean filterExists = !filterMap.isEmpty() &&
- filterMap.get(channel).stream().anyMatch(f ->
filter.id().equals(f.id()));
+ Set<PrioritizedFilter> filters = filterMap.computeIfAbsent(channel,
+ c -> new
ConcurrentSkipListSet<>(DynamicRouterConstants.FILTER_COMPARATOR));
+ List<PrioritizedFilterStatistics> filterStatistics =
filterStatisticsMap.computeIfAbsent(channel,
+ c -> Collections.synchronizedList(new ArrayList<>()));
+ boolean filterExists = filters.stream().anyMatch(f ->
filter.id().equals(f.id()));
boolean okToAdd = update == filterExists;
if (okToAdd) {
- Set<PrioritizedFilter> filters = filterMap.computeIfAbsent(channel,
- c -> new
ConcurrentSkipListSet<>(DynamicRouterConstants.FILTER_COMPARATOR));
- filters.add(filter);
- List<PrioritizedFilterStatistics> filterStatistics =
filterStatisticsMap.computeIfAbsent(channel,
- c -> Collections.synchronizedList(new ArrayList<>()));
+ // the set is ordered by priority and id: adding the updated
filter would neither replace a filter with
+ // the same priority nor remove the one with the old priority, so
the existing filter must be removed
+ // (its statistics stay, as when a filter is removed: they
represent actions that happened)
+ if (filterExists
+ && filters.stream().anyMatch(f ->
filter.id().equals(f.id()) && f.priority() == filter.priority())) {
+ // same priority: the set has no atomic replace, so remove the
existing filter first
+ filters.removeIf(f -> filter.id().equals(f.id()));
+ filters.add(filter);
+ } else {
+ // add the new filter first and then remove the old instance,
so that an exchange routed meanwhile
+ // still finds a filter for this subscription
+ filters.add(filter);
+ if (filterExists) {
+ filters.removeIf(f -> f != filter &&
filter.id().equals(f.id()));
+ }
+ }
filterStatistics.add(filter.statistics());
- LOG.debug("Added subscription: {}", filter);
+ LOG.debug("{} subscription: {}", filterExists ? "Updated" :
"Added", filter);
return filter.id();
}
return String.format("Error: Filter could not be %s -- existing filter
found with matching ID: %b",
diff --git
a/components/camel-dynamic-router/src/test/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterServiceUpdateTest.java
b/components/camel-dynamic-router/src/test/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterServiceUpdateTest.java
new file mode 100644
index 000000000000..2be4ed3f6730
--- /dev/null
+++
b/components/camel-dynamic-router/src/test/java/org/apache/camel/component/dynamicrouter/filter/DynamicRouterFilterServiceUpdateTest.java
@@ -0,0 +1,91 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.dynamicrouter.filter;
+
+import org.apache.camel.builder.PredicateBuilder;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+/**
+ * Updating a subscription replaces it, whatever its priority, with the filter
instances the service creates.
+ */
+class DynamicRouterFilterServiceUpdateTest {
+
+ static final String CHANNEL = "test";
+
+ DynamicRouterFilterService filterService;
+
+ @BeforeEach
+ void setup() {
+ filterService = new DynamicRouterFilterService();
+ filterService.initializeChannelFilters(CHANNEL);
+ filterService.addFilterForChannel("sub", 1,
PredicateBuilder.constant(true), "mock:old", CHANNEL, false);
+ }
+
+ @Test
+ void testUpdateWithSamePriorityReplacesTheSubscription() {
+ String result
+ = filterService.addFilterForChannel("sub", 1,
PredicateBuilder.constant(true), "mock:new", CHANNEL, true);
+
+ assertEquals("sub", result);
+ assertEquals(1, filterService.getFiltersForChannel(CHANNEL).size());
+ assertEquals("mock:new", filterService.getFilterById("sub",
CHANNEL).endpoint());
+ // the statistics of the replaced filter stay, as when a filter is
removed
+ assertEquals(2, filterService.getStatisticsForChannel(CHANNEL).size());
+ }
+
+ @Test
+ void testUpdateWithOtherPriorityReplacesTheSubscription() {
+ String result
+ = filterService.addFilterForChannel("sub", 10,
PredicateBuilder.constant(true), "mock:new", CHANNEL, true);
+
+ assertEquals("sub", result);
+ assertEquals(1, filterService.getFiltersForChannel(CHANNEL).size());
+ PrioritizedFilter filter = filterService.getFilterById("sub", CHANNEL);
+ assertEquals(10, filter.priority());
+ assertEquals("mock:new", filter.endpoint());
+ // the statistics of the replaced filter stay, as when a filter is
removed
+ assertEquals(2, filterService.getStatisticsForChannel(CHANNEL).size());
+ }
+
+ @Test
+ void testUpdateWithOtherPriorityKeepsTheNewInstance() {
+ // a lower and then a higher priority than the existing filter
+ for (int priority : new int[] { 0, 5 }) {
+ PrioritizedFilter updated = filterService.createFilter("sub",
priority, PredicateBuilder.constant(true),
+ "mock:new" + priority, new
PrioritizedFilterStatistics("sub"));
+
+ String result = filterService.addFilterForChannel(updated,
CHANNEL, true);
+
+ assertEquals("sub", result);
+ assertEquals(1,
filterService.getFiltersForChannel(CHANNEL).size());
+ assertSame(updated, filterService.getFilterById("sub", CHANNEL));
+ }
+ }
+
+ @Test
+ void testSubscribeToAChannelWithoutFilters() {
+ String result
+ = filterService.addFilterForChannel("other", 1,
PredicateBuilder.constant(true), "mock:other", "other", false);
+
+ assertEquals("other", result);
+ assertEquals(1, filterService.getFiltersForChannel("other").size());
+ }
+}