From f18b82a3a8a1b4170939639d6acab5c8141b0c70 Mon Sep 17 00:00:00 2001 From: smjain <49463903+allthingssecurity@users.noreply.github.com> Date: Thu, 24 Sep 2026 22:25:47 +0530 Subject: [PATCH 1/2] CAMEL-25012: camel-core - OnCompletion EIP: make a graceful shutdown wait for the parallel onCompletion tasks With parallelProcessing, OnCompletionProcessor submits the onCompletion of an exchange to its thread pool when the exchange's unit of work is done, and the exchange then leaves the inflight repository. The graceful shutdown waits for the route's inflight exchanges and for the pending exchanges of its ShutdownAware services, but OnCompletionProcessor was not ShutdownAware. So the shutdown did not wait for the onCompletion tasks, and then OnCompletionProcessor shut its thread pool down with shutdownNow: queued onCompletion tasks were dropped and running ones were interrupted, without any timeout being reported. OnCompletionProcessor is now ShutdownAware, and counts its onCompletion tasks as pending from when they are submitted until they are done, as CAMEL-24995 does for the Wire Tap EIP. The shutdownNow of the pool then only affects tasks that are still pending after the shutdown timeout. Co-Authored-By: Claude Opus 5.5 --- .../processor/OnCompletionProcessor.java | 62 ++++++++++--- ...pletionParallelProcessingShutdownTest.java | 90 +++++++++++++++++++ 2 files changed, 139 insertions(+), 13 deletions(-) create mode 100644 core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java index 07219abcc5563..1c707a2aa3dee 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java @@ -19,6 +19,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutorService; +import java.util.concurrent.atomic.LongAdder; import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; @@ -30,9 +31,11 @@ import org.apache.camel.Predicate; import org.apache.camel.Processor; import org.apache.camel.Route; +import org.apache.camel.ShutdownRunningTask; import org.apache.camel.Traceable; import org.apache.camel.spi.IdAware; import org.apache.camel.spi.RouteIdAware; +import org.apache.camel.spi.ShutdownAware; import org.apache.camel.spi.StepIdAware; import org.apache.camel.spi.SynchronizationRouteAware; import org.apache.camel.support.ExchangeHelper; @@ -47,7 +50,8 @@ /** * Processor implementing onCompletion. */ -public class OnCompletionProcessor extends BaseProcessorSupport implements Traceable, IdAware, RouteIdAware, StepIdAware { +public class OnCompletionProcessor extends BaseProcessorSupport + implements Traceable, ShutdownAware, IdAware, RouteIdAware, StepIdAware { private static final Logger LOG = LoggerFactory.getLogger(OnCompletionProcessor.class); @@ -64,6 +68,7 @@ public class OnCompletionProcessor extends BaseProcessorSupport implements Trace private final boolean useOriginalBody; private final boolean afterConsumer; private final boolean routeScoped; + private final LongAdder taskCount = new LongAdder(); public OnCompletionProcessor(CamelContext camelContext, Processor processor, ExecutorService executorService, boolean shutdownExecutorService, @@ -115,6 +120,22 @@ public CamelContext getCamelContext() { return camelContext; } + @Override + public boolean deferShutdown(ShutdownRunningTask shutdownRunningTask) { + // not in use + return true; + } + + @Override + public int getPendingExchangesSize() { + return taskCount.intValue(); + } + + @Override + public void prepareShutdown(boolean suspendOnly, boolean forced) { + // noop + } + @Override public String getId() { return id; @@ -162,6 +183,30 @@ public boolean process(Exchange exchange, AsyncCallback callback) { return true; } + /** + * Submits the onCompletion task to the thread pool (parallel processing). The task is counted as pending from when + * it is submitted until it is done, so a graceful shutdown waits for it. + */ + @SuppressWarnings("deprecation") + private void submitTask(Runnable task) { + taskCount.increment(); + Runnable counted = () -> { + try { + task.run(); + } finally { + taskCount.decrement(); + } + }; + try { + // Deprecated since 4.19.0 + executorService.submit(prepareMDCParallelTask(camelContext, counted)); + } catch (RuntimeException e) { + // the task will not run + taskCount.decrement(); + throw e; + } + } + protected boolean isCreateCopy() { // we need to create a correlated copy if we run in parallel mode or is in after consumer mode (as the UoW would be done on the original exchange otherwise) return executorService != null || afterConsumer; @@ -301,7 +346,6 @@ public void onAfterRoute(Route route, Exchange exchange) { }; } - @SuppressWarnings("deprecation") @Override public void onComplete(final Exchange exchange) { if (shouldSkip(exchange, onFailureOnly)) { @@ -316,9 +360,7 @@ public void onComplete(final Exchange exchange) { LOG.debug("Processing onComplete: {}", copy); doProcess(processor, copy); }; - // Deprecated since 4.19.0 - task = prepareMDCParallelTask(camelContext, task); - executorService.submit(task); + submitTask(task); } else { // run without thread-pool LOG.debug("Processing onComplete: {}", copy); @@ -326,7 +368,6 @@ public void onComplete(final Exchange exchange) { } } - @SuppressWarnings("deprecation") @Override public void onFailure(final Exchange exchange) { if (shouldSkip(exchange, onCompleteOnly)) { @@ -349,9 +390,7 @@ public void onFailure(final Exchange exchange) { // restore exception after processing copy.setException(original); }; - // Deprecated since 4.19.0 - task = prepareMDCParallelTask(camelContext, task); - executorService.submit(task); + submitTask(task); } else { // run without thread-pool LOG.debug("Processing onFailure: {}", copy); @@ -444,7 +483,6 @@ public void onBeforeRoute(Route route, Exchange exchange) { // NO-OP } - @SuppressWarnings("deprecation") @Override public void onAfterRoute(Route route, Exchange exchange) { LOG.debug("onAfterRoute from Route {}", route.getRouteId()); @@ -479,9 +517,7 @@ public void onAfterRoute(Route route, Exchange exchange) { LOG.debug("Processing onAfterRoute: {}", copy); doProcess(processor, copy); }; - // Deprecated since 4.19.0 - task = prepareMDCParallelTask(camelContext, task); - executorService.submit(task); + submitTask(task); } else { // run without thread-pool LOG.debug("Processing onAfterRoute: {}", copy); diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java new file mode 100644 index 0000000000000..3f69bd85ac21e --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java @@ -0,0 +1,90 @@ +/* + * 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.processor; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.ContextTestSupport; +import org.apache.camel.builder.RouteBuilder; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * The onCompletion tasks running in the thread pool of a parallel onCompletion are pending exchanges, so a graceful + * shutdown waits for them. + */ +class OnCompletionParallelProcessingShutdownTest extends ContextTestSupport { + + private final CountDownLatch onCompletionStarted = new CountDownLatch(1); + private final CountDownLatch camelStopping = new CountDownLatch(1); + private final AtomicInteger onCompletionDone = new AtomicInteger(); + private final ExecutorService stopper = Executors.newSingleThreadExecutor(); + + @AfterEach + void shutdownStopper() { + stopper.shutdownNow(); + } + + @Test + void testGracefulShutdownWaitsForOnCompletion() throws Exception { + template.sendBody("direct:start", "Hello World"); + assertTrue(onCompletionStarted.await(10, TimeUnit.SECONDS)); + OnCompletionProcessor onCompletion = context.getProcessor("oc", OnCompletionProcessor.class); + assertEquals(1, onCompletion.getPendingExchangesSize(), "The running onCompletion task should be pending"); + + // the exchange is done, and Camel is stopped while its onCompletion is still running + Future stop = stopper.submit(() -> { + context.stop(); + return null; + }); + await().atMost(10, TimeUnit.SECONDS).until(() -> context.isStopping() || context.isStopped()); + camelStopping.countDown(); + stop.get(30, TimeUnit.SECONDS); + + assertEquals(1, onCompletionDone.get(), "The onCompletion should be done"); + assertEquals(0, onCompletion.getPendingExchangesSize(), "No onCompletion task should be pending"); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("direct:start") + .onCompletion().id("oc").parallelProcessing() + .process(e -> { + onCompletionStarted.countDown(); + // the onCompletion takes until Camel is stopping + if (camelStopping.await(10, TimeUnit.SECONDS)) { + onCompletionDone.incrementAndGet(); + } + }) + .end() + .to("mock:result"); + } + }; + } +} From 2e6598a95133b6ff4820d00e959bb309faae3732 Mon Sep 17 00:00:00 2001 From: smjain <49463903+allthingssecurity@users.noreply.github.com> Date: Sun, 27 Sep 2026 16:51:19 +0530 Subject: [PATCH 2/2] CAMEL-25012: camel-core - OnCompletion EIP: do not count the tasks dropped by a forced shutdown as pending When a graceful shutdown times out, doShutdown calls shutdownNow on the thread pool of a parallel onCompletion. The tasks still queued in the pool never run, so they were never uncounted and stayed as pending exchanges. They are now subtracted using the list that shutdownNow returns. Running tasks are interrupted and uncount themselves as before. The counter is not reset when the processor is started: after a forced stopRoute the thread pool is not shut down, and its tasks keep running and uncount themselves after a startRoute, so a reset would make the counter negative. A restart of CamelContext creates new processors. Adds OnCompletionParallelProcessingForcedShutdownTest, and a 4.23 upgrade guide note that stopping or suspending a route now waits for the parallel onCompletion tasks. Co-Authored-By: Claude Opus 5.5 --- .../processor/OnCompletionProcessor.java | 6 +- ...nParallelProcessingForcedShutdownTest.java | 87 +++++++++++++++++++ .../pages/camel-4x-upgrade-guide-4_23.adoc | 11 +++ 3 files changed, 103 insertions(+), 1 deletion(-) create mode 100644 core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java index 1c707a2aa3dee..f6dd28ede3966 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java @@ -112,7 +112,11 @@ protected void doStop() throws Exception { protected void doShutdown() throws Exception { ServiceHelper.stopAndShutdownService(processor); if (shutdownExecutorService) { - getCamelContext().getExecutorServiceManager().shutdownNow(executorService); + List dropped = getCamelContext().getExecutorServiceManager().shutdownNow(executorService); + if (dropped != null && !dropped.isEmpty()) { + // the tasks still queued in the thread pool will never run, so they are no longer pending + taskCount.add(-dropped.size()); + } } } diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java new file mode 100644 index 0000000000000..3b6926d8093ee --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java @@ -0,0 +1,87 @@ +/* + * 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.processor; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.CamelContext; +import org.apache.camel.ContextTestSupport; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.builder.ThreadPoolProfileBuilder; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * When a graceful shutdown times out, the thread pool of a parallel onCompletion is shut down, which drops the + * onCompletion tasks still queued in it. They must no longer be counted as pending exchanges. + */ +class OnCompletionParallelProcessingForcedShutdownTest extends ContextTestSupport { + + private final CountDownLatch release = new CountDownLatch(1); + private final AtomicInteger started = new AtomicInteger(); + + @Test + void testForcedShutdownDropsQueuedOnCompletion() throws Exception { + // the first onCompletion runs (and waits), the second is queued in the pool with one thread + template.sendBody("direct:start", "A"); + template.sendBody("direct:start", "B"); + await().atMost(10, TimeUnit.SECONDS).until(() -> started.get() == 1); + OnCompletionProcessor onCompletion = context.getProcessor("oc", OnCompletionProcessor.class); + assertEquals(2, onCompletion.getPendingExchangesSize()); + + // the graceful shutdown times out, and the pool is shut down: the running task is interrupted, and the queued + // task is dropped + context.getShutdownStrategy().setTimeout(1); + context.stop(); + assertTrue(context.getShutdownStrategy().hasTimeoutOccurred()); + + await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(0, onCompletion.getPendingExchangesSize(), + "No onCompletion task should be pending after the thread pool is shut down")); + assertEquals(1, started.get(), "The queued onCompletion should not have run"); + } + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext context = super.createCamelContext(); + context.getExecutorServiceManager() + .registerThreadPoolProfile(new ThreadPoolProfileBuilder("oneThread").poolSize(1).maxPoolSize(1).build()); + return context; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("direct:start") + .onCompletion().id("oc").parallelProcessing().executorService("oneThread") + .process(e -> { + started.incrementAndGet(); + release.await(20, TimeUnit.SECONDS); + }) + .end() + .to("mock:result"); + } + }; + } +} diff --git a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc index f70b258eedc10..9f98bbeba24e5 100644 --- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc +++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc @@ -782,6 +782,17 @@ seen, and a repository backed by a remote store keeps its connection while the r forget the ids, clear the repository with `IdempotentRepository.clear()` (or the `clear` JMX operation of the Idempotent Consumer). +=== camel-core - a graceful shutdown waits for the parallel onCompletion tasks + +With `onCompletion().parallelProcessing()`, stopping or suspending a route (and stopping `CamelContext`) now waits +for the onCompletion tasks that are running or queued in its thread pool, up to the shutdown timeout, just as it +waits for the inflight exchanges. Previously the shutdown did not wait for them, and when the thread pool was shut +down, the queued tasks were dropped and the running ones were interrupted. + +An onCompletion with `parallelProcessing` that synchronously stops its own route now waits for itself until the +shutdown timeout occurs, and the route is then stopped forcibly. Stop the route asynchronously instead, for example +from a separate thread or with the Control Bus `async=true` option. + === Component deprecation ==== camel-minio