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 38581a08b446 CAMEL-25274: camel-dapr - acknowledge a pub/sub event 
when the exchange is done, and ask for a redelivery when it failed (#27303)
38581a08b446 is described below

commit 38581a08b4465111060e19e1f26848ef82d26142
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 02:07:30 2026 +0530

    CAMEL-25274: camel-dapr - acknowledge a pub/sub event when the exchange is 
done, and ask for a redelivery when it failed (#27303)
    
    `DaprPubSubConsumer`'s listener returned `Mono.just(Status.SUCCESS)` as 
soon as the exchange was handed to the route. That status is the 
acknowledgement of the event in the Dapr streaming subscription (the SDK sends 
it to the sidecar as the `TopicEventResponse` status), so an event whose route 
failed was acknowledged and dropped by Dapr, and an event was acknowledged 
before an asynchronous route had processed it.
    This change: the listener returns a `Mono` that emits when the exchange is 
done: `SUCCESS` when it completed, `RETRY` when it failed (the status the SDK 
itself uses when a listener fails). The upgrade guide for 4.23 gets a note, 
since events of failed exchanges are now redelivered.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../dapr/consumer/DaprPubSubConsumer.java          | 25 +++++--
 .../dapr/consumer/DaprPubSubConsumerTest.java      | 85 +++++++++++++++++++++-
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    | 10 +++
 3 files changed, 112 insertions(+), 8 deletions(-)

diff --git 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
index b156dd7f7015..5448b99b3631 100644
--- 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
+++ 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
@@ -119,13 +119,24 @@ public class DaprPubSubConsumer extends DefaultConsumer {
 
         @Override
         public Mono<Status> onEvent(CloudEvent<byte[]> cloudEvent) {
-            final Exchange exchange = createServiceBusExchange(cloudEvent);
-
-            // use default consumer callback
-            AsyncCallback cb = defaultConsumerCallback(exchange, true);
-            getAsyncProcessor().process(exchange, cb);
-
-            return Mono.just(Status.SUCCESS);
+            // the status is the acknowledgement of the event: answer it when 
the exchange is done, and
+            // ask Dapr to redeliver the event when the exchange failed (or 
was marked rollback only) instead of
+            // dropping it
+            return Mono.create(sink -> {
+                final Exchange exchange = createServiceBusExchange(cloudEvent);
+
+                // use default consumer callback
+                AsyncCallback cb = defaultConsumerCallback(exchange, true);
+                getAsyncProcessor().process(exchange, doneSync -> {
+                    Status status = exchange.isFailed() || 
exchange.isRollbackOnly() || exchange.isRollbackOnlyLast()
+                            ? Status.RETRY : Status.SUCCESS;
+                    try {
+                        cb.done(doneSync);
+                    } finally {
+                        sink.success(status);
+                    }
+                });
+            });
         }
 
         @Override
diff --git 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
index da7bee99b824..9461469ec887 100644
--- 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
+++ 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
@@ -18,6 +18,7 @@ package org.apache.camel.component.dapr.consumer;
 
 import java.nio.charset.StandardCharsets;
 import java.time.OffsetDateTime;
+import java.util.concurrent.CompletableFuture;
 
 import io.dapr.client.DaprClientBuilder;
 import io.dapr.client.DaprPreviewClient;
@@ -40,14 +41,17 @@ import org.apache.camel.test.junit6.CamelTestSupport;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.mockito.ArgumentCaptor;
+import reactor.core.publisher.Mono;
 
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyBoolean;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
@@ -125,7 +129,12 @@ public class DaprPubSubConsumerTest extends 
CamelTestSupport {
         when(cloudEvent.getTraceParent()).thenReturn(traceParent);
         when(cloudEvent.getTraceState()).thenReturn(traceState);
 
-        listenerCaptor.getValue().onEvent(cloudEvent).block();
+        doAnswer(inv -> {
+            inv.getArgument(1, AsyncCallback.class).done(true);
+            return true;
+        }).when(processor).process(any(Exchange.class), 
any(AsyncCallback.class));
+
+        assertEquals(SubscriptionListener.Status.SUCCESS, 
listenerCaptor.getValue().onEvent(cloudEvent).block());
 
         verify(processor).process(exchangeCaptor.capture(), 
callbackCaptor.capture());
 
@@ -148,4 +157,78 @@ public class DaprPubSubConsumerTest extends 
CamelTestSupport {
         verify(mockSubscription).close();
         verify(mockClient).close();
     }
+
+    @Test
+    void testFailedExchangeIsRedelivered() throws Exception {
+        consumer.doStart();
+
+        doAnswer(inv -> {
+            inv.getArgument(0, Exchange.class).setException(new 
IllegalStateException("Forced"));
+            inv.getArgument(1, AsyncCallback.class).done(true);
+            return true;
+        }).when(processor).process(any(Exchange.class), 
any(AsyncCallback.class));
+
+        SubscriptionListener.Status status = 
listenerCaptor.getValue().onEvent(newCloudEvent()).block();
+
+        // a failed exchange must not acknowledge the event, Dapr would drop it
+        assertEquals(SubscriptionListener.Status.RETRY, status);
+    }
+
+    @Test
+    void testRollbackOnlyExchangeIsRedelivered() throws Exception {
+        consumer.doStart();
+
+        doAnswer(inv -> {
+            // a rollback without an exception, as markRollbackOnly() does
+            inv.getArgument(0, Exchange.class).setRollbackOnly(true);
+            inv.getArgument(1, AsyncCallback.class).done(true);
+            return true;
+        }).when(processor).process(any(Exchange.class), 
any(AsyncCallback.class));
+
+        SubscriptionListener.Status status = 
listenerCaptor.getValue().onEvent(newCloudEvent()).block();
+
+        assertEquals(SubscriptionListener.Status.RETRY, status);
+    }
+
+    @Test
+    void testRollbackOnlyLastExchangeIsRedelivered() throws Exception {
+        consumer.doStart();
+
+        doAnswer(inv -> {
+            // a rollback without an exception, as markRollbackOnlyLast() does
+            inv.getArgument(0, Exchange.class).setRollbackOnlyLast(true);
+            inv.getArgument(1, AsyncCallback.class).done(true);
+            return true;
+        }).when(processor).process(any(Exchange.class), 
any(AsyncCallback.class));
+
+        SubscriptionListener.Status status = 
listenerCaptor.getValue().onEvent(newCloudEvent()).block();
+
+        assertEquals(SubscriptionListener.Status.RETRY, status);
+    }
+
+    @Test
+    void testEventIsAcknowledgedWhenTheExchangeIsDone() throws Exception {
+        consumer.doStart();
+
+        // the route continues asynchronously: the callback is called later
+        doReturn(false).when(processor).process(any(Exchange.class), 
any(AsyncCallback.class));
+
+        Mono<SubscriptionListener.Status> result = 
listenerCaptor.getValue().onEvent(newCloudEvent());
+        CompletableFuture<SubscriptionListener.Status> status = 
result.toFuture();
+
+        verify(processor).process(exchangeCaptor.capture(), 
callbackCaptor.capture());
+        assertFalse(status.isDone(), "The event must not be acknowledged 
before the exchange is done");
+
+        exchangeCaptor.getValue().setException(new 
IllegalStateException("Forced"));
+        callbackCaptor.getValue().done(false);
+
+        assertEquals(SubscriptionListener.Status.RETRY, status.get());
+    }
+
+    private static CloudEvent<byte[]> newCloudEvent() {
+        CloudEvent<byte[]> cloudEvent = mock(CloudEvent.class);
+        when(cloudEvent.getData()).thenReturn("testBody".getBytes());
+        when(cloudEvent.getId()).thenReturn("myId");
+        return cloudEvent;
+    }
 }
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 c1a368e51a05..858917b437dd 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
@@ -792,6 +792,16 @@ the maximum reconsume times of the consumer group, 16 by 
default, then it goes t
 A route that fails for a message, and relied on the message being dropped, now 
receives it again; handle the
 exception in the route (for example with `onException(...).handled(true)`) to 
acknowledge the message anyway.
 
+=== camel-dapr - pub/sub consumer acknowledges an event after the exchange
+
+The pub/sub consumer acknowledged every event with `SUCCESS` as soon as the 
exchange was handed to the route, also
+when the route then failed, so Dapr dropped events whose processing failed. 
The consumer now answers when the exchange
+is done: `SUCCESS` when it completed, and `RETRY` when it failed (or was 
marked rollback only), so that Dapr redelivers
+the event. A route that fails for an event, and relied on the event being 
dropped, now receives it again; handle the
+exception in the route (for example with `onException(...).handled(true)`) to 
acknowledge the event anyway. A poison
+event that always fails now gets `RETRY` every time, until a Dapr resiliency 
policy for the pub/sub component or a
+dead-letter topic on the subscription limits the redeliveries.
+
 === Components and Language removal
 
 ==== camel-csimple, camel-csimple-joor and csimple-maven-plugin

Reply via email to