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 7c6191be6856 CAMEL-25003: camel-core - Do not hand out or stop in-use
evicted non-singleton producers in ServicePool
7c6191be6856 is described below
commit 7c6191be68562c1e47066263f5f3f12c808d95d1
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 22:28:05 2026 +0530
CAMEL-25003: camel-core - Do not hand out or stop in-use evicted
non-singleton producers in ServicePool
For endpoints whose producer is not a singleton (ftp/sftp/ftps, smb, ssh,
and the polling consumers of pollEnrich), ServicePool keeps a queue of idle
producers per endpoint and an LRU cache of the producers it handed out.
When the LRU evicted such a producer it was stopped on the next acquire or
release but not taken out of the queue, so a stopped producer could be
handed out again, and a producer in use when evicted was stopped under the
running exchange. This is a regression from CAMEL-20829 (4.7).
MultiplePool now tracks its non-evicted producers: evict takes an idle
producer out of the queue before stopping it, release stops a producer
that was evicted while in use instead of offering it back, and stop() also
stops pending evicts. onEvict no longer stops a non-singleton producer
whose pool is gone, since it can only be in use and release stops it. A
pool being stopped marks itself stopped, so a producer released or evicted
into it meanwhile is stopped rather than left behind.
Closes #26861
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.../ProducerCacheNonSingletonEvictionTest.java | 216 +++++++++++++++++++++
.../apache/camel/support/cache/ServicePool.java | 35 +++-
2 files changed, 246 insertions(+), 5 deletions(-)
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/ProducerCacheNonSingletonEvictionTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/ProducerCacheNonSingletonEvictionTest.java
new file mode 100644
index 000000000000..5976c5b5fc7d
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/ProducerCacheNonSingletonEvictionTest.java
@@ -0,0 +1,216 @@
+/*
+ * 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.util.Map;
+
+import org.apache.camel.AsyncCallback;
+import org.apache.camel.Component;
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.support.DefaultAsyncProducer;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.cache.DefaultProducerCache;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Non-singleton producers are pooled per endpoint. When one of them is
evicted from the producer cache, it must not be
+ * handed out again once it is stopped, and it must not be stopped while it is
in use.
+ */
+class ProducerCacheNonSingletonEvictionTest extends ContextTestSupport {
+
+ private DefaultProducerCache cache;
+ private Endpoint a;
+ private Endpoint b;
+ private Endpoint c;
+ private Endpoint d;
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Override
+ @BeforeEach
+ public void setUp() throws Exception {
+ super.setUp();
+ context.addComponent("pooled", new PooledComponent());
+ a = context.getEndpoint("pooled:a");
+ b = context.getEndpoint("pooled:b");
+ c = context.getEndpoint("pooled:c");
+ d = context.getEndpoint("pooled:d");
+ cache = new DefaultProducerCache(this, context, 2);
+ cache.start();
+ }
+
+ @AfterEach
+ void stopCache() {
+ cache.stop();
+ }
+
+ @Test
+ void testEvictIdleProducer() {
+ // two producers for endpoint a, both idle in its pool
+ PooledProducer a1 = acquire(a);
+ PooledProducer a2 = acquire(a);
+ cache.releaseProducer(a, a1);
+ cache.releaseProducer(a, a2);
+
+ // a producer for endpoint b evicts a1 (the eldest) from the cache,
and a1 is stopped
+ cache.releaseProducer(b, acquire(b));
+ cache.cleanUp();
+ assertTrue(a1.isStopped(), "Evicted idle producer should be stopped");
+
+ // the stopped producer must not be handed out again
+ PooledProducer next = acquire(a);
+ assertTrue(next.isStarted(), "Acquired producer should be started");
+ assertSame(a2, next);
+ cache.releaseProducer(a, next);
+ }
+
+ @Test
+ void testEvictProducerInUse() {
+ // a1 is in use, a2 is idle in the pool of endpoint a
+ PooledProducer a1 = acquire(a);
+ cache.releaseProducer(a, acquire(a));
+
+ // a producer for endpoint b evicts a1 (the eldest) from the cache
while it is in use
+ cache.releaseProducer(b, acquire(b));
+ cache.cleanUp();
+ cache.releaseProducer(a, acquire(a));
+ assertTrue(a1.isStarted(), "Evicted producer should not be stopped
while it is in use");
+
+ // when it is released, it is stopped instead of being returned to the
pool
+ cache.releaseProducer(a, a1);
+ assertTrue(a1.isStopped(), "Evicted producer should be stopped when it
is released");
+
+ PooledProducer next1 = acquire(a);
+ PooledProducer next2 = acquire(a);
+ assertNotSame(a1, next1);
+ assertNotSame(a1, next2);
+ assertTrue(next1.isStarted(), "Acquired producer should be started");
+ assertTrue(next2.isStarted(), "Acquired producer should be started");
+ cache.releaseProducer(a, next1);
+ cache.releaseProducer(a, next2);
+ }
+
+ @Test
+ void testStopPoolWhileProducersInUse() {
+ // b1 is idle in the pool of endpoint b; a1 and a2 are in use and
evict b1 from the cache
+ cache.releaseProducer(b, acquire(b));
+ PooledProducer a1 = acquire(a);
+ PooledProducer a2 = acquire(a);
+ cache.cleanUp();
+
+ // c1 evicts a1, and with three pools for a cache size of 2 the pool
of endpoint a is stopped
+ cache.releaseProducer(c, acquire(c));
+ cache.cleanUp();
+ assertTrue(a1.isStarted(), "Producer should not be stopped while it is
in use");
+ assertTrue(a2.isStarted(), "Producer should not be stopped while it is
in use");
+
+ // d1 evicts a2, whose pool no longer exists
+ cache.releaseProducer(d, acquire(d));
+ cache.cleanUp();
+ assertTrue(a2.isStarted(), "Producer should not be stopped while it is
in use");
+
+ // when they are released, they are stopped instead of being pooled
again
+ cache.releaseProducer(a, a1);
+ cache.releaseProducer(a, a2);
+ assertTrue(a1.isStopped(), "Producer of a stopped pool should be
stopped when it is released");
+ assertTrue(a2.isStopped(), "Producer of a stopped pool should be
stopped when it is released");
+
+ PooledProducer next = acquire(a);
+ assertNotSame(a1, next);
+ assertNotSame(a2, next);
+ assertTrue(next.isStarted(), "Acquired producer should be started");
+ cache.releaseProducer(a, next);
+ }
+
+ @Test
+ void testStopPoolWithEvictedIdleProducer() {
+ // b1 is idle in the pool of endpoint b, a1 is in use
+ PooledProducer b1 = acquire(b);
+ cache.releaseProducer(b, b1);
+ PooledProducer a1 = acquire(a);
+
+ // c1 evicts b1, and with three pools for a cache size of 2 the pool
of endpoint b is stopped,
+ // which must stop the evicted idle producer as the pool is no longer
cleaned up afterwards
+ cache.releaseProducer(c, acquire(c));
+ cache.cleanUp();
+ assertTrue(b1.isStopped(), "Evicted idle producer of a stopped pool
should be stopped");
+ assertTrue(a1.isStarted(), "Producer should not be stopped while it is
in use");
+ cache.releaseProducer(a, a1);
+ }
+
+ private PooledProducer acquire(Endpoint endpoint) {
+ return (PooledProducer) cache.acquireProducer(endpoint);
+ }
+
+ private static final class PooledComponent extends DefaultComponent {
+
+ @Override
+ protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
+ return new PooledEndpoint(uri, this);
+ }
+ }
+
+ private static final class PooledEndpoint extends DefaultEndpoint {
+
+ PooledEndpoint(String uri, Component component) {
+ super(uri, component);
+ }
+
+ @Override
+ public Producer createProducer() {
+ return new PooledProducer(this);
+ }
+
+ @Override
+ public Consumer createConsumer(Processor processor) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public boolean isSingletonProducer() {
+ return false;
+ }
+ }
+
+ private static final class PooledProducer extends DefaultAsyncProducer {
+
+ PooledProducer(Endpoint endpoint) {
+ super(endpoint);
+ }
+
+ @Override
+ public boolean process(Exchange exchange, AsyncCallback callback) {
+ callback.done(true);
+ return true;
+ }
+ }
+}
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 626c685ddaed..d6e61b9f60a6 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
@@ -19,6 +19,7 @@ package org.apache.camel.support.cache;
import java.util.ArrayList;
import java.util.Deque;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
@@ -94,9 +95,14 @@ abstract class ServicePool<S extends Service> extends
ServiceSupport implements
// the pool is growing too large, so we need to stop (stop
will remove itself from pool)
p.stop();
}
- } else {
+ } else if (e.isSingletonProducer()) {
// service no longer in a pool (such as being released twice, or
can happen during shutdown of Camel etc)
stopAndRemove(s);
+ } else {
+ // no pool for this endpoint (for example it was stopped, which
stopped its idle services, or the
+ // endpoint of the service is not the instance used as the pool
key): do not stop the service here, as it
+ // may be in use or idle in a live pool; a service in use is
stopped when it is released without a pool
+ LOG.trace("Evicted service: {} is no longer in a pool", s);
}
}
@@ -324,11 +330,15 @@ abstract class ServicePool<S extends Service> extends
ServiceSupport implements
private final Endpoint endpoint;
private final BlockingQueue<S> queue;
private final Deque<S> evicts;
+ // the services created by this pool which have not been evicted, only
these are kept in the queue
+ private final Set<S> active;
+ private volatile boolean stopped;
MultiplePool(Endpoint endpoint) {
this.endpoint = endpoint;
this.queue = new ArrayBlockingQueue<>(capacity);
this.evicts = new ConcurrentLinkedDeque<>();
+ this.active = ConcurrentHashMap.newKeySet();
}
private void cleanupEvicts() {
@@ -345,6 +355,7 @@ abstract class ServicePool<S extends Service> extends
ServiceSupport implements
if (s == null) {
s = creator.apply(endpoint);
s.start();
+ active.add(s);
}
return s;
}
@@ -353,8 +364,13 @@ abstract class ServicePool<S extends Service> extends
ServiceSupport implements
public void release(S s) {
cleanupEvicts();
- if (!queue.offer(s)) {
- // there is no room so let's just stop and discard this
+ if (stopped || !active.contains(s) || !queue.offer(s)) {
+ // the pool is stopped, it was evicted while in use, or there
is no room so let's just stop and discard this
+ active.remove(s);
+ doStop(s);
+ } else if ((stopped || !active.contains(s)) && queue.remove(s)) {
+ // the pool was stopped, or the service evicted, after the
check above, and did not take it from the queue
+ active.remove(s);
doStop(s);
}
}
@@ -366,16 +382,25 @@ abstract class ServicePool<S extends Service> extends
ServiceSupport implements
@Override
public void stop() {
+ stopped = true;
ArrayList<S> list = new ArrayList<>();
queue.drainTo(list);
pool.remove(endpoint);
list.forEach(this::doStop);
+ cleanupEvicts();
}
@Override
public void evict(S s) {
- // to be evicted
- evicts.add(s);
+ // only an idle service can be stopped (by cleanupEvicts): take it
out of the queue so it is not acquired
+ // again, and a service in use is stopped when it is released
+ if (active.remove(s) && queue.remove(s)) {
+ evicts.add(s);
+ if (stopped) {
+ // the pool was stopped meanwhile, and may have stopped
its evicts already
+ cleanupEvicts();
+ }
+ }
}
@Override