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");
+ }
+ };
+ }
+}