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 323e370e67d1 CAMEL-24933, CAMEL-24934, CAMEL-16073: camel-pulsar - 
consumer acknowledgement and exchange handling
323e370e67d1 is described below

commit 323e370e67d12eca1142bf5d5458a0e85c2700df
Author: Andrea Cosentino <[email protected]>
AuthorDate: Mon Sep 28 18:45:37 2026 +0200

    CAMEL-24933, CAMEL-24934, CAMEL-16073: camel-pulsar - consumer 
acknowledgement and exchange handling
    
    CAMEL-24933: PulsarMessageUtils.updateExchange returned a copy of the
    exchange, so the pooled exchange the consumer created was never released,
    one per message. It now populates and returns the exchange it was given.
    
    CAMEL-24934: a failed acknowledgement is now reported with its cause.
    
    CAMEL-16073: when the route fails and allowManualAcknowledgement is false
    (the default), the consumer negatively acknowledges the message instead of
    leaving it to the acknowledgement timeout. It passes the Message, so a
    configured negativeAckRedeliveryBackoff sees the redelivery count and
    escalates. The nack is sent before the exception handler runs. Redelivery
    now follows negativeAckRedeliveryDelayMicros (60 seconds by default); the
    upgrade guide shows how to keep the previous 10 second timing.
    
    Closes #26779
    
    Co-authored-by: Claude Opus 5 <[email protected]>
---
 .../component/pulsar/PulsarMessageListener.java    |  21 ++-
 .../pulsar/utils/message/PulsarMessageUtils.java   |  15 +-
 .../PulsarMessageListenerAcknowledgementTest.java  | 197 +++++++++++++++++++++
 .../utils/message/PulsarMessageUtilsTest.java      |  34 ++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  24 +++
 5 files changed, 282 insertions(+), 9 deletions(-)

diff --git 
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
 
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
index 9038c354d06d..a3481bf770e7 100644
--- 
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
+++ 
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
@@ -50,14 +50,16 @@ public class PulsarMessageListener implements 
MessageListener<byte[]> {
         pulsarConsumer.getAsyncProcessor().process(exchange, doneSync -> {
             try {
                 if (exchange.getException() != null) {
+                    // tell the broker first: a custom ExceptionHandler that 
throws must not cost us the
+                    // negative acknowledgement, which would silently fall 
back to ack-timeout redelivery
+                    negativeAcknowledge(consumer, message);
                     
pulsarConsumer.getExceptionHandler().handleException("Error processing 
exchange", exchange,
                             exchange.getException());
                 } else {
                     try {
                         acknowledge(consumer, message);
                     } catch (Exception e) {
-                        
pulsarConsumer.getExceptionHandler().handleException("Error processing 
exchange", exchange,
-                                exchange.getException());
+                        
pulsarConsumer.getExceptionHandler().handleException("Error acknowledging 
message", exchange, e);
                     }
                 }
             } finally {
@@ -73,4 +75,19 @@ public class PulsarMessageListener implements 
MessageListener<byte[]> {
         }
     }
 
+    /**
+     * Tells the broker the message was not processed, so that it is 
redelivered after
+     * {@code negativeAckRedeliveryDelayMicros} instead of waiting for the 
acknowledgement timeout. Left to the route
+     * when manual acknowledgement is enabled, the same way {@link 
#acknowledge} is.
+     * <p>
+     * The message is passed rather than its id on purpose: only that overload 
carries the redelivery count into
+     * {@code NegativeAcksTracker}, and without it a configured {@code 
negativeAckRedeliveryBackoff} is always asked for
+     * the delay of attempt zero and never escalates.
+     */
+    private void negativeAcknowledge(final Consumer<byte[]> consumer, final 
Message<byte[]> message) {
+        if (!endpoint.getPulsarConfiguration().isAllowManualAcknowledgement()) 
{
+            consumer.negativeAcknowledge(message);
+        }
+    }
+
 }
diff --git 
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
 
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
index d03b3fe1287c..a9ff18556622 100644
--- 
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
+++ 
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
@@ -41,10 +41,13 @@ public final class PulsarMessageUtils {
     private PulsarMessageUtils() {
     }
 
-    public static Exchange updateExchange(final Message<byte[]> message, final 
Exchange input) {
-        final Exchange output = input.copy();
-
-        org.apache.camel.Message msg = output.getIn();
+    /**
+     * Populates the given exchange from the Pulsar message and returns it. 
The exchange is updated in place: copying it
+     * would orphan the exchange the consumer took from the exchange factory, 
which with a pooled factory is never
+     * returned to the pool.
+     */
+    public static Exchange updateExchange(final Message<byte[]> message, final 
Exchange exchange) {
+        org.apache.camel.Message msg = exchange.getIn();
 
         msg.setHeader(EVENT_TIME, message.getEventTime());
         msg.setHeader(MESSAGE_ID, message.getMessageId());
@@ -60,9 +63,7 @@ public final class PulsarMessageUtils {
 
         msg.setBody(message.getValue());
 
-        output.setIn(msg);
-
-        return output;
+        return exchange;
     }
 
     public static Exchange updateExchangeWithException(final Exception 
exception, final Exchange input) {
diff --git 
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarMessageListenerAcknowledgementTest.java
 
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarMessageListenerAcknowledgementTest.java
new file mode 100644
index 000000000000..85261d991b10
--- /dev/null
+++ 
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarMessageListenerAcknowledgementTest.java
@@ -0,0 +1,197 @@
+/*
+ * 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.pulsar;
+
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.ExceptionHandler;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.pulsar.client.api.Consumer;
+import org.apache.pulsar.client.api.ConsumerBuilder;
+import org.apache.pulsar.client.api.Message;
+import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Answers.RETURNS_SELF;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * A route failure used to leave the message neither acknowledged nor 
negatively acknowledged, so it was only
+ * redelivered once the acknowledgement timeout expired (CAMEL-16073), and a 
failing acknowledgement was reported
+ * without its cause.
+ */
+public class PulsarMessageListenerAcknowledgementTest extends CamelTestSupport 
{
+
+    private static final String ENDPOINT_URI
+            = 
"pulsar:persistent://public/default/camel-ack-test?subscriptionName=camel-subscription";
+
+    private final MessageId messageId = mock(MessageId.class);
+
+    @Test
+    public void testSuccessfulExchangeIsAcknowledged() throws Exception {
+        final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+
+        listener().received(pulsarConsumer, message("ok"));
+
+        verify(pulsarConsumer, timeout(5000)).acknowledge(messageId);
+        verify(pulsarConsumer, never()).negativeAcknowledge(messageId);
+    }
+
+    @Test
+    public void testFailedExchangeIsNegativelyAcknowledged() throws Exception {
+        final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+        final Message<byte[]> message = message("fail");
+
+        listener().received(pulsarConsumer, message);
+
+        // the Message overload, not the MessageId one: only that carries the 
redelivery count into
+        // NegativeAcksTracker, so a configured negativeAckRedeliveryBackoff 
can escalate
+        verify(pulsarConsumer, timeout(5000)).negativeAcknowledge(message);
+        verify(pulsarConsumer, never()).negativeAcknowledge(messageId);
+        verify(pulsarConsumer, never()).acknowledge(messageId);
+    }
+
+    @Test
+    public void testManualAcknowledgementLeavesTheFailedMessageToTheRoute() 
throws Exception {
+        final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+        final Message<byte[]> message = message("fail");
+
+        final CapturingExceptionHandler exceptionHandler = new 
CapturingExceptionHandler();
+        pulsarConsumer().setExceptionHandler(exceptionHandler);
+
+        context.getEndpoint(ENDPOINT_URI, 
PulsarEndpoint.class).getPulsarConfiguration()
+                .setAllowManualAcknowledgement(true);
+        try {
+            listener().received(pulsarConsumer, message);
+
+            // wait for the callback to have run before asserting that nothing 
was sent to the broker
+            assertTrue(exceptionHandler.latch.await(5, TimeUnit.SECONDS),
+                    "the route failure should still be reported");
+            verify(pulsarConsumer, never()).negativeAcknowledge(message);
+            verify(pulsarConsumer, never()).negativeAcknowledge(messageId);
+            verify(pulsarConsumer, never()).acknowledge(messageId);
+        } finally {
+            context.getEndpoint(ENDPOINT_URI, 
PulsarEndpoint.class).getPulsarConfiguration()
+                    .setAllowManualAcknowledgement(false);
+        }
+    }
+
+    @Test
+    public void testAcknowledgeFailureIsReportedWithItsCause() throws 
Exception {
+        final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+        final PulsarClientException failure = new 
PulsarClientException("cannot acknowledge");
+        doThrow(failure).when(pulsarConsumer).acknowledge(messageId);
+
+        final CapturingExceptionHandler exceptionHandler = new 
CapturingExceptionHandler();
+        pulsarConsumer().setExceptionHandler(exceptionHandler);
+
+        listener().received(pulsarConsumer, message("ok"));
+
+        assertTrue(exceptionHandler.latch.await(5, TimeUnit.SECONDS), "the 
acknowledgement failure should be reported");
+        assertNotNull(exceptionHandler.captured.get(), "the acknowledgement 
failure should carry its cause");
+        assertEquals("cannot acknowledge", 
exceptionHandler.captured.get().getMessage());
+    }
+
+    private PulsarMessageListener listener() {
+        return new PulsarMessageListener(context.getEndpoint(ENDPOINT_URI, 
PulsarEndpoint.class), pulsarConsumer());
+    }
+
+    private PulsarConsumer pulsarConsumer() {
+        return (PulsarConsumer) context.getRoutes().get(0).getConsumer();
+    }
+
+    private Message<byte[]> message(String body) {
+        final Message<byte[]> message = mock(Message.class);
+        
when(message.getValue()).thenReturn(body.getBytes(StandardCharsets.UTF_8));
+        when(message.getMessageId()).thenReturn(messageId);
+        when(message.getProperties()).thenReturn(Collections.emptyMap());
+        return message;
+    }
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        final CamelContext context = super.createCamelContext();
+
+        final ConsumerBuilder<byte[]> builder = mock(ConsumerBuilder.class, 
RETURNS_SELF);
+        when(builder.subscribe()).thenReturn(mock(Consumer.class));
+
+        final PulsarClient pulsarClient = mock(PulsarClient.class);
+        when(pulsarClient.newConsumer()).thenReturn(builder);
+
+        final PulsarComponent component = new PulsarComponent(context);
+        component.setPulsarClient(pulsarClient);
+        context.addComponent("pulsar", component);
+
+        return context;
+    }
+
+    @Override
+    protected RoutesBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from(ENDPOINT_URI)
+                        .process(exchange -> {
+                            if 
("fail".equals(exchange.getIn().getBody(String.class))) {
+                                throw new IllegalStateException("simulated 
route failure");
+                            }
+                        })
+                        .to("mock:result");
+            }
+        };
+    }
+
+    private static final class CapturingExceptionHandler implements 
ExceptionHandler {
+
+        private final CountDownLatch latch = new CountDownLatch(1);
+        private final AtomicReference<Throwable> captured = new 
AtomicReference<>();
+
+        @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) {
+            captured.compareAndSet(null, exception);
+            latch.countDown();
+        }
+    }
+}
diff --git 
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
 
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
index 18deba91db0f..1e98017524ce 100644
--- 
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
+++ 
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
@@ -17,10 +17,20 @@
 package org.apache.camel.component.pulsar.utils.message;
 
 import java.io.Serializable;
+import java.util.Collections;
 
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.pulsar.client.api.Message;
+import org.apache.pulsar.client.api.MessageId;
 import org.junit.jupiter.api.Test;
 
+import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
 
 public class PulsarMessageUtilsTest {
 
@@ -48,6 +58,30 @@ public class PulsarMessageUtilsTest {
 
         assertNotNull(expected);
     }
+
+    @Test
+    public void testUpdateExchangeUpdatesTheGivenExchangeInPlace() {
+        final DefaultCamelContext context = new DefaultCamelContext();
+        final Exchange exchange = new DefaultExchange(context);
+
+        final MessageId messageId = mock(MessageId.class);
+        final Message<byte[]> message = mock(Message.class);
+        when(message.getValue()).thenReturn("Hello World!".getBytes());
+        when(message.getMessageId()).thenReturn(messageId);
+        when(message.getKey()).thenReturn("aKey");
+        when(message.getTopicName()).thenReturn("aTopic");
+        when(message.getProducerName()).thenReturn("aProducer");
+        when(message.getProperties()).thenReturn(Collections.emptyMap());
+
+        final Exchange updated = PulsarMessageUtils.updateExchange(message, 
exchange);
+
+        // a copy here would orphan the exchange the consumer took from the 
exchange factory
+        assertSame(exchange, updated);
+        assertEquals(messageId, 
updated.getIn().getHeader(PulsarMessageHeaders.MESSAGE_ID));
+        assertEquals("aKey", 
updated.getIn().getHeader(PulsarMessageHeaders.KEY));
+        assertEquals("aTopic", 
updated.getIn().getHeader(PulsarMessageHeaders.TOPIC_NAME));
+        assertEquals("Hello World!", updated.getIn().getBody(String.class));
+    }
 }
 
 class Obj implements Serializable {
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 7afd41e043c9..22e3d3ae84a2 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
@@ -3192,6 +3192,30 @@ reply of a discarded exchange (`waitForTaskToComplete`) 
is released with that ex
 `timeout`, or forever when the timeout is disabled. On completions handed over 
to a discarded InOnly exchange, such as
 the commit or rollback of the consumer that received the message, now run as a 
failure, where previously they never ran.
 
+=== camel-pulsar - PulsarMessageUtils.updateExchange returns the exchange it 
was given
+
+`PulsarMessageUtils.updateExchange(message, exchange)` used to return a *copy* 
of the exchange passed to
+it. It now populates and returns that same instance.
+
+The copy orphaned the exchange the consumer had taken from the exchange 
factory, so with
+`camel.main.exchange-factory=pooled` every consumed message leaked one pooled 
exchange and the pool never
+refilled. Code outside the component that called this method and relied on 
getting an independent copy
+must make its own copy instead.
+
+=== camel-pulsar - a failed exchange is negatively acknowledged
+
+When a route fails, the consumer now calls `negativeAcknowledge` on the Pulsar 
consumer instead of
+leaving the message unacknowledged. This only applies when 
`allowManualAcknowledgement` is `false`
+(the default); with manual acknowledgement the route stays in charge, as 
before.
+
+This changes when the message comes back. Previously it was redelivered once 
the acknowledgement
+timeout expired, which `camel-pulsar` sets to 10 seconds by default through 
`ackTimeoutMillis`. A
+negative acknowledgement removes the message from the client's 
unacknowledged-message tracker, so
+redelivery now follows `negativeAckRedeliveryDelayMicros`, which defaults to 
60 seconds, and honours
+`negativeAckRedeliveryBackoff` when one is configured.
+
+A route that wants the previous timing can set 
`negativeAckRedeliveryDelayMicros=10000000`.
+
 === camel-seda - multipleConsumers broadcasts to consumers with different uri 
options
 
 With `multipleConsumers=true` every consumer of a SEDA queue now receives a 
copy of each message, also when the consumers

Reply via email to