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());
+    }
+}

Reply via email to