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 c8432fed6faa CAMEL-25012: camel-core - OnCompletion EIP: make a 
graceful shutdown wait for the parallel onCompletion tasks (#26870)
c8432fed6faa is described below

commit c8432fed6faa93c966f58be2a405ff951f5a64c2
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 17:53:18 2026 +0530

    CAMEL-25012: camel-core - OnCompletion EIP: make a graceful shutdown wait 
for the parallel onCompletion tasks (#26870)
    
    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. When a graceful shutdown times
    out, the tasks dropped by shutdownNow are subtracted from the pending
    count so they do not stay counted.
    
    Adds 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 <[email protected]>
---
 .../camel/processor/OnCompletionProcessor.java     | 68 ++++++++++++----
 ...letionParallelProcessingForcedShutdownTest.java | 87 +++++++++++++++++++++
 ...OnCompletionParallelProcessingShutdownTest.java | 90 ++++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    | 11 +++
 4 files changed, 242 insertions(+), 14 deletions(-)

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 07219abcc556..f6dd28ede396 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 @@ package org.apache.camel.processor;
 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.Ordered;
 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 @@ import static org.apache.camel.util.ObjectHelper.notNull;
 /**
  * Processor implementing <a 
href="http://camel.apache.org/oncompletion.html";>onCompletion</a>.
  */
-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,
@@ -107,7 +112,11 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
     protected void doShutdown() throws Exception {
         ServiceHelper.stopAndShutdownService(processor);
         if (shutdownExecutorService) {
-            
getCamelContext().getExecutorServiceManager().shutdownNow(executorService);
+            List<Runnable> 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());
+            }
         }
     }
 
@@ -115,6 +124,22 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
         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 +187,30 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
         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 +350,6 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
             };
         }
 
-        @SuppressWarnings("deprecation")
         @Override
         public void onComplete(final Exchange exchange) {
             if (shouldSkip(exchange, onFailureOnly)) {
@@ -316,9 +364,7 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
                     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 +372,6 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
             }
         }
 
-        @SuppressWarnings("deprecation")
         @Override
         public void onFailure(final Exchange exchange) {
             if (shouldSkip(exchange, onCompleteOnly)) {
@@ -349,9 +394,7 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
                     // 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 +487,6 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
                     // NO-OP
                 }
 
-                @SuppressWarnings("deprecation")
                 @Override
                 public void onAfterRoute(Route route, Exchange exchange) {
                     LOG.debug("onAfterRoute from Route {}", 
route.getRouteId());
@@ -479,9 +521,7 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport implements Trace
                             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/OnCompletionParallelProcessingForcedShutdownTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java
new file mode 100644
index 000000000000..3b6926d8093e
--- /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/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 000000000000..3f69bd85ac21
--- /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");
+            }
+        };
+    }
+}
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 ea0f34c2c218..a85c1efe1c2f 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
@@ -940,6 +940,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

Reply via email to