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