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 b3c138ab811b CAMEL-25004: camel-core - Producer cache and stream
caching: fix lifecycle bugs found in a deep review (#26862)
b3c138ab811b is described below
commit b3c138ab811b1fdc8500edef84afccbdcdbed0cf
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 09:05:11 2026 +0200
CAMEL-25004: camel-core - Producer cache and stream caching: fix lifecycle
bugs found in a deep review (#26862)
- 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..626c685ddaed 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;
@@ -185,6 +187,27 @@ abstract class ServicePool<S extends Service> extends
ServiceSupport implements
singlePoolEvicted.clear();
}
+ /**
+ * 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;
+ }
+
/**
* Stops the service safely
*/
@@ -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);
}