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 4236c7fe66cb CAMEL-25012: camel-core - Wire Tap EIP: do not count the 
tapped exchanges dropped by a forced shutdown as pending (#26992)
4236c7fe66cb is described below

commit 4236c7fe66cbdf3d644db1a235995d34b99538e7
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 19:16:36 2026 +0530

    CAMEL-25012: camel-core - Wire Tap EIP: do not count the tapped exchanges 
dropped by a forced shutdown as pending (#26992)
    
    When a graceful shutdown times out, doShutdown calls shutdownNow on the
    thread pool of the wire tap. The tapped exchanges still queued in the pool
    are never sent, so their tasks never uncounted themselves and stayed as
    pending exchanges. They are now subtracted using the list that shutdownNow
    returns. Running tasks are interrupted and uncount themselves as before.
    
    This is the same fix as for the parallel onCompletion in CAMEL-25012.
    
    Adds WireTapForcedShutdownTest.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../apache/camel/processor/WireTapProcessor.java   |  7 +-
 .../camel/processor/WireTapForcedShutdownTest.java | 90 ++++++++++++++++++++++
 2 files changed, 96 insertions(+), 1 deletion(-)

diff --git 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/WireTapProcessor.java
 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/WireTapProcessor.java
index e7fa9e38bd22..9f41f0fd3f17 100644
--- 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/WireTapProcessor.java
+++ 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/WireTapProcessor.java
@@ -17,6 +17,7 @@
 package org.apache.camel.processor;
 
 import java.io.IOException;
+import java.util.List;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.atomic.LongAdder;
 
@@ -380,7 +381,11 @@ public class WireTapProcessor extends BaseProcessorSupport
     protected void doShutdown() throws Exception {
         ServiceHelper.stopAndShutdownServices(processorExchangeFactory, 
taskFactory, 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());
+            }
         }
     }
 }
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/WireTapForcedShutdownTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/WireTapForcedShutdownTest.java
new file mode 100644
index 000000000000..0e867c9aa126
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/WireTapForcedShutdownTest.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.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 the wire tap is shut 
down, which drops the tapped exchanges
+ * still queued in it. They must no longer be counted as pending exchanges.
+ */
+class WireTapForcedShutdownTest extends ContextTestSupport {
+
+    private final CountDownLatch release = new CountDownLatch(1);
+    private final AtomicInteger started = new AtomicInteger();
+
+    @Test
+    void testForcedShutdownDropsQueuedTap() throws Exception {
+        try {
+            // the first tapped exchange 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);
+            WireTapProcessor tap = context.getProcessor("tap", 
WireTapProcessor.class);
+            assertEquals(2, tap.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, 
tap.getPendingExchangesSize(),
+                            "No tapped exchange should be pending after the 
thread pool is shut down"));
+            assertEquals(1, started.get(), "The queued tapped exchange should 
not have been sent");
+        } finally {
+            release.countDown();
+        }
+    }
+
+    @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").wireTap("direct:tap").executorService("oneThread").id("tap").to("mock:result");
+
+                from("direct:tap")
+                        .process(e -> {
+                            started.incrementAndGet();
+                            release.await(20, TimeUnit.SECONDS);
+                        });
+            }
+        };
+    }
+}

Reply via email to