This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch fix/CAMEL-25004
in repository https://gitbox.apache.org/repos/asf/camel.git

commit f0e7a8f1d0bce00d47b8e7d73a6047055c371f53
Author: Claus Ibsen <[email protected]>
AuthorDate: Thu Sep 24 19:13:08 2026 +0200

    CAMEL-25004: camel-core - Producer cache and stream caching: fix lifecycle 
bugs found in a deep review
    
    - Evicting the producer of a singleton endpoint stopped the endpoint
      also when the routes used it. It is now only stopped when not in use
      (a dynamic endpoint that no route consumes from).
    - The stream caching strategy added its threshold spool rules, the
      classes given by name and the core converters again each time it was
      started. It no longer duplicates them.
    - allowClasses/denyClasses given as both classes and names failed with
      UnsupportedOperationException.
    - The spool directory was not removed when spooling only by used heap
      memory (or custom spool rules).
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../impl/engine/DefaultStreamCachingStrategy.java  | 51 +++++++-----
 .../impl/ProducerCacheEvictEndpointInUseTest.java  | 69 ++++++++++++++++
 .../impl/StreamCachingStrategyRestartTest.java     | 94 ++++++++++++++++++++++
 .../apache/camel/support/cache/ServicePool.java    | 28 ++++++-
 4 files changed, 222 insertions(+), 20 deletions(-)

diff --git 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultStreamCachingStrategy.java
 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultStreamCachingStrategy.java
index b0c93e6bbd2b..7a43f3e3c90c 100644
--- 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultStreamCachingStrategy.java
+++ 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultStreamCachingStrategy.java
@@ -72,6 +72,10 @@ public class DefaultStreamCachingStrategy extends 
ServiceSupport implements Came
     private boolean removeSpoolDirectoryWhenStopping = true;
     private final UtilizationStatistics statistics = new 
UtilizationStatistics();
     private final Set<SpoolRule> spoolRules = new LinkedHashSet<>();
+    // the spool rules added when starting (and removed when stopping), from 
the spool thresholds
+    private final List<SpoolRule> thresholdSpoolRules = new ArrayList<>();
+    // whether spooling to disk is in use (the spool directory is only removed 
when it was in use)
+    private volatile boolean spoolInUse;
     private volatile boolean anySpoolRules;
 
     @Override
@@ -377,28 +381,16 @@ public class DefaultStreamCachingStrategy extends 
ServiceSupport implements Came
         }
 
         // find core type converters that can convert to StreamCache
+        // (clear first, as the strategy may be started again after a restart)
+        coreConverters.clear();
         var set = 
getCamelContext().getTypeConverterRegistry().lookup(StreamCache.class).entrySet();
         set.forEach(e -> coreConverters.add(new CoreConverter(e.getKey(), 
e.getValue())));
 
         if (allowClassNames != null) {
-            if (allowClasses == null) {
-                allowClasses = new ArrayList<>();
-            }
-            for (String name : allowClassNames.split(",")) {
-                name = name.trim();
-                Class<?> clazz = 
camelContext.getClassResolver().resolveMandatoryClass(name);
-                allowClasses.add(clazz);
-            }
+            allowClasses = resolveClasses(allowClasses, allowClassNames);
         }
         if (denyClassNames != null) {
-            if (denyClasses == null) {
-                denyClasses = new ArrayList<>();
-            }
-            for (String name : denyClassNames.split(",")) {
-                name = name.trim();
-                Class<?> clazz = 
camelContext.getClassResolver().resolveMandatoryClass(name);
-                denyClasses.add(clazz);
-            }
+            denyClasses = resolveClasses(denyClasses, denyClassNames);
         }
 
         if (spoolUsedHeapMemoryThreshold > 99) {
@@ -439,15 +431,17 @@ public class DefaultStreamCachingStrategy extends 
ServiceSupport implements Came
                 }
             }
             if (spoolThreshold > 0) {
-                spoolRules.add(new FixedThresholdSpoolRule());
+                thresholdSpoolRules.add(new FixedThresholdSpoolRule());
             }
             if (spoolUsedHeapMemoryThreshold > 0) {
                 if (spoolUsedHeapMemoryLimit == null) {
                     // use max by default
                     spoolUsedHeapMemoryLimit = SpoolUsedHeapMemoryLimit.Max;
                 }
-                spoolRules.add(new 
UsedHeapMemorySpoolRule(spoolUsedHeapMemoryLimit));
+                thresholdSpoolRules.add(new 
UsedHeapMemorySpoolRule(spoolUsedHeapMemoryLimit));
             }
+            spoolRules.addAll(thresholdSpoolRules);
+            spoolInUse = true;
         }
 
         LOG.debug("StreamCaching configuration {}", this);
@@ -468,6 +462,11 @@ public class DefaultStreamCachingStrategy extends 
ServiceSupport implements Came
             LOG.debug("Removing spool directory: {}", spoolDirectory);
             FileUtil.removeDir(spoolDirectory);
         }
+        spoolInUse = false;
+
+        // remove the spool rules added when starting, as they are added again 
if started again
+        thresholdSpoolRules.forEach(spoolRules::remove);
+        thresholdSpoolRules.clear();
 
         if (LOG.isDebugEnabled() && statistics.isStatisticsEnabled()) {
             LOG.debug("Stopping StreamCachingStrategy with statistics: {}", 
statistics);
@@ -476,8 +475,22 @@ public class DefaultStreamCachingStrategy extends 
ServiceSupport implements Came
         statistics.reset();
     }
 
+    private Collection<Class<?>> resolveClasses(Collection<Class<?>> classes, 
String names) throws ClassNotFoundException {
+        // use a new list, as the existing may be immutable (such as set via 
setAllowClasses)
+        Collection<Class<?>> answer = classes != null ? new 
ArrayList<>(classes) : new ArrayList<>();
+        for (String name : names.split(",")) {
+            Class<?> clazz = 
camelContext.getClassResolver().resolveMandatoryClass(name.trim());
+            // avoid duplicates when started again after a restart
+            if (!answer.contains(clazz)) {
+                answer.add(clazz);
+            }
+        }
+        return answer;
+    }
+
     private boolean isSpoolRemovable() {
-        return spoolThreshold > 0 && spoolDirectory != null && 
isRemoveSpoolDirectoryWhenStopping();
+        // spooling may be in use by any of the spool rules (not only the 
spool threshold)
+        return spoolInUse && spoolDirectory != null && 
isRemoveSpoolDirectoryWhenStopping();
     }
 
     @Override
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/impl/ProducerCacheEvictEndpointInUseTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/impl/ProducerCacheEvictEndpointInUseTest.java
new file mode 100644
index 000000000000..8a2d266db8b1
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/impl/ProducerCacheEvictEndpointInUseTest.java
@@ -0,0 +1,69 @@
+/*
+ * 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.impl;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.service.ServiceHelper;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Evicting the producer of a singleton endpoint from a producer cache must 
not stop the endpoint when it is in use by
+ * the routes, but stops a dynamic endpoint that is not in use.
+ */
+public class ProducerCacheEvictEndpointInUseTest extends ContextTestSupport {
+
+    @Test
+    public void testEvictDoesNotStopEndpointInUse() throws Exception {
+        template.sendBodyAndHeader("direct:start", "A", "uri", "seda:foo");
+        // the cache holds one producer, so this evicts the producer of 
seda:foo
+        template.sendBodyAndHeader("direct:start", "B", "uri", "seda:bar");
+        // the eviction is cleaned up on the next use of the cache
+        template.sendBodyAndHeader("direct:start", "C", "uri", "seda:bar");
+
+        Endpoint foo = context.hasEndpoint("seda:foo");
+        assertTrue(ServiceHelper.isStarted(foo), "seda:foo is used by a route 
and must not be stopped");
+    }
+
+    @Test
+    public void testEvictStopsDynamicEndpointNotInUse() throws Exception {
+        // seda:dynamic is only used by toD, so it is a dynamic endpoint
+        template.sendBodyAndHeader("direct:start", "A", "uri", "seda:dynamic");
+        template.sendBodyAndHeader("direct:start", "B", "uri", "seda:bar");
+        template.sendBodyAndHeader("direct:start", "C", "uri", "seda:bar");
+
+        Endpoint dynamic = context.hasEndpoint("seda:dynamic");
+        assertFalse(ServiceHelper.isStarted(dynamic), "the dynamic endpoint 
not in use should be stopped to free resources");
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:start").toD("${header.uri}", 1);
+
+                from("seda:foo").to("mock:foo");
+                from("seda:bar").to("mock:bar");
+            }
+        };
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/impl/StreamCachingStrategyRestartTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/impl/StreamCachingStrategyRestartTest.java
new file mode 100644
index 000000000000..3013826aebe6
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/impl/StreamCachingStrategyRestartTest.java
@@ -0,0 +1,94 @@
+/*
+ * 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.impl;
+
+import java.io.File;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.spi.StreamCachingStrategy;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The stream caching strategy can be stopped and started again without 
accumulating its configuration, and it removes
+ * its spool directory also when spooling is only by used heap memory.
+ */
+public class StreamCachingStrategyRestartTest extends ContextTestSupport {
+
+    @Override
+    public boolean isUseRouteBuilder() {
+        return false;
+    }
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        CamelContext context = super.createCamelContext();
+        context.setStreamCaching(true);
+        return context;
+    }
+
+    @Test
+    public void testRestartDoesNotDuplicateAllowClasses() throws Exception {
+        StreamCachingStrategy strategy = context.getStreamCachingStrategy();
+        strategy.stop();
+        strategy.setEnabled(true);
+        strategy.setAllowClasses("java.lang.String");
+
+        strategy.start();
+        assertEquals(1, strategy.getAllowClasses().size());
+
+        // stop and start the strategy again
+        strategy.stop();
+        strategy.start();
+        assertEquals(1, strategy.getAllowClasses().size(), "the allow classes 
should not be added again on restart");
+    }
+
+    @Test
+    public void testAllowClassesAndAllowClassNames() throws Exception {
+        StreamCachingStrategy strategy = context.getStreamCachingStrategy();
+        context.stop();
+        // classes and class names can be combined
+        strategy.setAllowClasses(Integer.class);
+        strategy.setAllowClasses("java.lang.String");
+
+        context.start();
+        assertEquals(2, strategy.getAllowClasses().size());
+    }
+
+    @Test
+    public void testSpoolDirectoryRemovedWhenSpoolingByHeapMemory() throws 
Exception {
+        File dir = testDirectory("spool").toFile();
+        CamelContext camel = new DefaultCamelContext();
+        camel.setStreamCaching(true);
+        StreamCachingStrategy strategy = camel.getStreamCachingStrategy();
+        strategy.setSpoolEnabled(true);
+        strategy.setSpoolDirectory(dir);
+        // spool only by used heap memory
+        strategy.setSpoolThreshold(0);
+        strategy.setSpoolUsedHeapMemoryThreshold(1);
+
+        camel.start();
+        assertTrue(dir.exists(), "the spool directory should be created");
+
+        camel.stop();
+        assertFalse(dir.exists(), "the spool directory should be removed when 
stopping");
+    }
+}
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java
 
b/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java
index 7dd2e62a2cbc..797df7aa8614 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java
@@ -26,8 +26,10 @@ import java.util.concurrent.ConcurrentLinkedDeque;
 import java.util.concurrent.ConcurrentMap;
 import java.util.function.Function;
 
+import org.apache.camel.CamelContext;
 import org.apache.camel.Endpoint;
 import org.apache.camel.NonManagedService;
+import org.apache.camel.Route;
 import org.apache.camel.Service;
 import org.apache.camel.support.LRUCache;
 import org.apache.camel.support.LRUCacheFactory;
@@ -188,6 +190,27 @@ abstract class ServicePool<S extends Service> extends 
ServiceSupport implements
     /**
      * Stops the service safely
      */
+    /**
+     * Whether the endpoint is (still) in use by the routes, and must 
therefore not be stopped when its producer is
+     * evicted: an endpoint that is static in the endpoint registry (resolved 
when the routes were setup), or that a
+     * route is consuming from.
+     */
+    private static boolean isEndpointInUse(Endpoint endpoint) {
+        CamelContext context = endpoint.getCamelContext();
+        if (context == null) {
+            return false;
+        }
+        if (context.getEndpointRegistry().isStatic(endpoint.getEndpointUri())) 
{
+            return true;
+        }
+        for (Route route : context.getRoutes()) {
+            if (route.getEndpoint() == endpoint) {
+                return true;
+            }
+        }
+        return false;
+    }
+
     private static <S extends Service> void stop(S s) {
         try {
             s.stop();
@@ -271,7 +294,10 @@ abstract class ServicePool<S extends Service> extends 
ServiceSupport implements
                 for (Map.Entry<Endpoint, Pool<S>> entry : 
singlePoolEvicted.entrySet()) {
                     Endpoint e = entry.getKey();
                     Pool<S> p = entry.getValue();
-                    doStop(e);
+                    if (!isEndpointInUse(e)) {
+                        // stop the endpoint as well (such as a dynamic 
endpoint from toD) to free its resources
+                        doStop(e);
+                    }
                     p.stop();
                     singlePoolEvicted.remove(e);
                 }

Reply via email to