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 e8cc7079fe9e CAMEL-24952: camel-seda - virtualThreadPerTask consumer 
stop waits for dispatched tasks that have not started
e8cc7079fe9e is described below

commit e8cc7079fe9edcbf67fa5814a884d8bf9f0e7709
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 22:25:26 2026 +0530

    CAMEL-24952: camel-seda - virtualThreadPerTask consumer stop waits for 
dispatched tasks that have not started
    
    ThreadPerTaskSedaConsumer incremented activeTasks only when a task
    started running. An exchange that was polled but whose task had not
    started yet was neither in the queue, inflight, nor counted, so a
    graceful stop considered the consumer idle. The task then ran against
    the stopped route and was rejected, and the message was lost.
    
    activeTasks is now incremented before taskExecutor.execute, and
    decremented, with the concurrency permit released, in the task's finally
    block and also when execute throws, where the permit previously leaked.
    
    Closes #26798
    
    Co-authored-by: Claude Opus 5.5 <[email protected]>
---
 .../component/seda/ThreadPerTaskSedaConsumer.java  |  20 ++-
 .../seda/ThreadPerTaskSedaConsumerStopTest.java    | 153 +++++++++++++++++++++
 2 files changed, 172 insertions(+), 1 deletion(-)

diff --git 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumer.java
 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumer.java
index a764c9e2fe39..3380efdd1755 100644
--- 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumer.java
+++ 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumer.java
@@ -128,9 +128,27 @@ public class ThreadPerTaskSedaConsumer extends 
SedaConsumer {
 
     @Override
     protected void processPolledExchange(Exchange exchange) {
+        // count the task when it is dispatched (not when it starts running), 
so a graceful shutdown
+        // also waits for polled exchanges whose task has not started yet
+        activeTasks.increment();
+        boolean dispatched = false;
+        try {
+            dispatch(exchange);
+            dispatched = true;
+        } finally {
+            if (!dispatched) {
+                // the task was not dispatched (e.g. rejected), so it will 
never run and undo the count itself
+                activeTasks.decrement();
+                if (concurrencyLimiter != null) {
+                    concurrencyLimiter.release();
+                }
+            }
+        }
+    }
+
+    private void dispatch(Exchange exchange) {
         // Dispatch to task executor for processing
         taskExecutor.execute(() -> {
-            activeTasks.increment();
             try {
                 // Prepare the exchange
                 Exchange prepared = prepareExchange(exchange);
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumerStopTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumerStopTest.java
new file mode 100644
index 000000000000..20510860be67
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/ThreadPerTaskSedaConsumerStopTest.java
@@ -0,0 +1,153 @@
+/*
+ * 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.component.seda;
+
+import java.util.List;
+import java.util.concurrent.AbstractExecutorService;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.support.DefaultThreadPoolFactory;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A graceful stop of a virtualThreadPerTask seda consumer must wait for an 
exchange that has been polled and
+ * dispatched, but whose task has not started running yet.
+ */
+class ThreadPerTaskSedaConsumerStopTest extends ContextTestSupport {
+
+    private final CountDownLatch dispatched = new CountDownLatch(1);
+    private final CountDownLatch startTask = new CountDownLatch(1);
+    private final CountDownLatch awaitingTermination = new CountDownLatch(1);
+
+    /**
+     * Task executor which delays the start of the tasks, until the test lets 
them start.
+     */
+    private final class DelayedStartExecutorService extends 
AbstractExecutorService {
+        private final ExecutorService delegate;
+
+        DelayedStartExecutorService(ExecutorService delegate) {
+            this.delegate = delegate;
+        }
+
+        @Override
+        public void execute(Runnable task) {
+            delegate.execute(() -> {
+                try {
+                    startTask.await(20, TimeUnit.SECONDS);
+                } catch (InterruptedException e) {
+                    Thread.currentThread().interrupt();
+                }
+                task.run();
+            });
+            dispatched.countDown();
+        }
+
+        @Override
+        public boolean awaitTermination(long timeout, TimeUnit unit) throws 
InterruptedException {
+            awaitingTermination.countDown();
+            return delegate.awaitTermination(timeout, unit);
+        }
+
+        @Override
+        public void shutdown() {
+            delegate.shutdown();
+        }
+
+        @Override
+        public List<Runnable> shutdownNow() {
+            return delegate.shutdownNow();
+        }
+
+        @Override
+        public boolean isShutdown() {
+            return delegate.isShutdown();
+        }
+
+        @Override
+        public boolean isTerminated() {
+            return delegate.isTerminated();
+        }
+    }
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        CamelContext context = super.createCamelContext();
+        context.getExecutorServiceManager().setThreadPoolFactory(new 
DefaultThreadPoolFactory() {
+            @Override
+            public ExecutorService newCachedThreadPool(ThreadFactory 
threadFactory) {
+                ExecutorService answer = 
super.newCachedThreadPool(threadFactory);
+                // the task executor of the thread-per-task seda consumer (its 
coordinator is a single thread executor)
+                if (threadFactory.newThread(() -> {
+                }).getName().endsWith("seda://v")) {
+                    answer = new DelayedStartExecutorService(answer);
+                }
+                return answer;
+            }
+        });
+        return context;
+    }
+
+    @Test
+    void testStopWaitsForDispatchedExchange() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:result");
+        mock.expectedBodiesReceived("Hello");
+
+        template.sendBody("seda:v", "Hello");
+        // the exchange has been polled and dispatched, but its task has not 
started
+        assertTrue(dispatched.await(10, TimeUnit.SECONDS));
+
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        try {
+            Future<?> stop = executor.submit(() -> {
+                context.getRouteController().stopRoute("v");
+                return null;
+            });
+            // let the task start when the stop waits for the tasks to 
complete (or when the stop completed)
+            await().atMost(20, TimeUnit.SECONDS).until(() -> stop.isDone() || 
awaitingTermination.getCount() == 0);
+            startTask.countDown();
+            stop.get(20, TimeUnit.SECONDS);
+        } finally {
+            startTask.countDown();
+            executor.shutdownNow();
+        }
+
+        // the exchange was processed by the route before it was stopped
+        mock.assertIsSatisfied();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("seda:v?virtualThreadPerTask=true&pollTimeout=100").routeId("v").to("mock:result");
+            }
+        };
+    }
+}

Reply via email to