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 17f4f62c7eb3 CAMEL-25037: camel-support - an interrupted 
EventDrivenPollingConsumer.receive() must return (#26912)
17f4f62c7eb3 is described below

commit 17f4f62c7eb36f88c47899d5208001ddd49e8651
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 17:56:46 2026 +0530

    CAMEL-25037: camel-support - an interrupted 
EventDrivenPollingConsumer.receive() must return (#26912)
    
    Cause: receive() loops while the consumer is running, and on an
    InterruptedException calls handleInterruptedException, which since
    CAMEL-20297 restores the interrupt status of the thread (and logs through
    the interrupted exception handler). The loop then waits on the queue again,
    which throws at once because the thread is still interrupted.
    
    Effect: a thread interrupted in receive() never leaves it, even when a
    message is queued. It uses a full core and logs one WARN per iteration
    (hundreds of thousands of lines per second). This happens to a
    ConsumerTemplate.receive() or a pollEnrich with the default timeout whose
    thread is cancelled (Future.cancel(true), shutdownNow() of its executor, or
    a forced shutdown of a thread pool by Camel).
    
    Fix: keep the interrupt status, as CAMEL-20297 intended, and return null
    from receive() instead of waiting again, like receive(timeout) already does.
    The javadoc of receive() documents when null is returned.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../EventDrivenPollingConsumerInterruptTest.java   | 120 +++++++++++++++++++++
 .../camel/support/EventDrivenPollingConsumer.java  |   8 ++
 2 files changed, 128 insertions(+)

diff --git 
a/core/camel-core/src/test/java/org/apache/camel/impl/EventDrivenPollingConsumerInterruptTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/impl/EventDrivenPollingConsumerInterruptTest.java
new file mode 100644
index 000000000000..508777be2e12
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/impl/EventDrivenPollingConsumerInterruptTest.java
@@ -0,0 +1,120 @@
+/*
+ * 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.impl;
+
+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.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.spi.ExceptionHandler;
+import org.apache.camel.support.EventDrivenPollingConsumer;
+import org.apache.camel.support.service.ServiceHelper;
+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.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A receive() whose thread is interrupted while it waits for a message must 
return, and keep the interrupt status.
+ */
+class EventDrivenPollingConsumerInterruptTest extends ContextTestSupport {
+
+    private final ExecutorService executor = 
Executors.newSingleThreadExecutor();
+    private EventDrivenPollingConsumer pollingConsumer;
+
+    @Override
+    @AfterEach
+    public void tearDown() throws Exception {
+        // a receive() that is still waiting ends once the polling consumer is 
stopped
+        ServiceHelper.stopAndShutdownService(pollingConsumer);
+        executor.shutdownNow();
+        super.tearDown();
+    }
+
+    @Override
+    public boolean isUseRouteBuilder() {
+        return false;
+    }
+
+    @Test
+    void testInterruptedReceiveReturns() throws Exception {
+        context.start();
+        pollingConsumer = assertIsInstanceOf(EventDrivenPollingConsumer.class,
+                context.getEndpoint("direct:idle").createPollingConsumer());
+        AtomicInteger interruptedExceptions = new AtomicInteger();
+        pollingConsumer.setInterruptedExceptionHandler(new ExceptionHandler() {
+            @Override
+            public void handleException(Throwable exception) {
+                handleException(null, null, exception);
+            }
+
+            @Override
+            public void handleException(String message, Throwable exception) {
+                handleException(message, null, exception);
+            }
+
+            @Override
+            public void handleException(String message, Exchange exchange, 
Throwable exception) {
+                if (exception instanceof InterruptedException) {
+                    interruptedExceptions.incrementAndGet();
+                }
+            }
+        });
+        pollingConsumer.start();
+
+        AtomicReference<Thread> receiver = new AtomicReference<>();
+        AtomicBoolean interruptedAfterReceive = new AtomicBoolean();
+        Future<Exchange> future = executor.submit(() -> {
+            receiver.set(Thread.currentThread());
+            Exchange answer = pollingConsumer.receive();
+            
interruptedAfterReceive.set(Thread.currentThread().isInterrupted());
+            return answer;
+        });
+        await().atMost(10, TimeUnit.SECONDS).until(() -> 
isWaitingInReceive(receiver.get()));
+
+        receiver.get().interrupt();
+
+        // receive() returns instead of waiting again (which fails at once as 
the thread is still interrupted)
+        await().atMost(5, TimeUnit.SECONDS).until(future::isDone);
+        assertNull(future.get());
+        assertTrue(interruptedAfterReceive.get(), "The interrupt status of the 
thread should be kept");
+        assertEquals(1, interruptedExceptions.get());
+    }
+
+    private static boolean isWaitingInReceive(Thread thread) {
+        if (thread == null
+                || thread.getState() != Thread.State.WAITING && 
thread.getState() != Thread.State.TIMED_WAITING) {
+            return false;
+        }
+        for (StackTraceElement element : thread.getStackTrace()) {
+            if 
(EventDrivenPollingConsumer.class.getName().equals(element.getClassName())
+                    && "receive".equals(element.getMethodName())) {
+                return true;
+            }
+        }
+        return false;
+    }
+}
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/support/EventDrivenPollingConsumer.java
 
b/core/camel-support/src/main/java/org/apache/camel/support/EventDrivenPollingConsumer.java
index 039cdf30faac..df587500df0b 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/support/EventDrivenPollingConsumer.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/support/EventDrivenPollingConsumer.java
@@ -133,6 +133,12 @@ public class EventDrivenPollingConsumer extends 
PollingConsumerSupport implement
         return receive(0);
     }
 
+    /**
+     * Waits until a message is available and then returns it.
+     * <p/>
+     * Returns <tt>null</tt> if this consumer is stopped while waiting, or if 
the calling thread is interrupted while
+     * waiting. In the latter case the interrupt status of the thread is kept.
+     */
     @Override
     public Exchange receive() {
         // must be started
@@ -150,7 +156,9 @@ public class EventDrivenPollingConsumer extends 
PollingConsumerSupport implement
                         return answer;
                     }
                 } catch (InterruptedException e) {
+                    // the interrupt status is kept, so waiting again would 
fail at once: stop waiting
                     handleInterruptedException(e);
+                    return null;
                 }
             }
         } finally {

Reply via email to