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 81177cc0f59e CAMEL-25250: camel-grpc - end a streaming call with an 
error when the exchange failed, with the AGGREGATION and PROPAGATION strategies 
(#27247)
81177cc0f59e is described below

commit 81177cc0f59ef3bbe0d815a8c59bb9def585db2c
Author: allthingssecurity <[email protected]>
AuthorDate: Fri Oct 2 13:23:12 2026 +0530

    CAMEL-25250: camel-grpc - end a streaming call with an error when the 
exchange failed, with the AGGREGATION and PROPAGATION strategies (#27247)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../component/grpc/server/GrpcMethodHandler.java   |  25 ++-
 .../server/GrpcRequestAbstractStreamObserver.java  |  13 ++
 .../GrpcRequestAggregationStreamObserver.java      |   4 +
 .../GrpcRequestPropagationStreamObserver.java      |  21 +-
 .../grpc/GrpcConsumerStreamingExceptionTest.java   | 224 +++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |   9 +
 6 files changed, 284 insertions(+), 12 deletions(-)

diff --git 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcMethodHandler.java
 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcMethodHandler.java
index 6b0ddb052f9b..da5b410d8616 100644
--- 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcMethodHandler.java
+++ 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcMethodHandler.java
@@ -21,6 +21,7 @@ import java.util.List;
 import java.util.Map;
 
 import io.grpc.Status;
+import io.grpc.StatusRuntimeException;
 import io.grpc.stub.StreamObserver;
 import org.apache.camel.Exchange;
 import org.apache.camel.component.grpc.GrpcConstants;
@@ -66,15 +67,7 @@ public class GrpcMethodHandler {
         invokeRoute(endpoint, exchange);
 
         if (exchange.isFailed()) {
-            // the description IS transmitted to the client, so it carries the 
route's exception message only when
-            // the endpoint opted out of muting; the cause below stays local 
either way
-            String description = endpoint.getConfiguration().isMuteException()
-                    ? MUTED_DESCRIPTION : exchange.getException().getMessage();
-            responseObserver.onError(Status.INTERNAL
-                    .withDescription(description)
-                    // This can be attached to the Status locally, but NOT 
transmitted to the client!
-                    .withCause(exchange.getException())
-                    .asRuntimeException());
+            responseObserver.onError(toStatusException(endpoint, 
exchange.getException()));
         } else {
             Object responseBody = exchange.getIn().getBody();
             if (responseBody instanceof List) {
@@ -87,6 +80,20 @@ public class GrpcMethodHandler {
         }
     }
 
+    /**
+     * The error to send to the client when the exchange failed.
+     */
+    static StatusRuntimeException toStatusException(GrpcEndpoint endpoint, 
Exception cause) {
+        // the description IS transmitted to the client, so it carries the 
route's exception message only when
+        // the endpoint opted out of muting; the cause below stays local 
either way
+        String description = endpoint.getConfiguration().isMuteException() ? 
MUTED_DESCRIPTION : cause.getMessage();
+        return Status.INTERNAL
+                .withDescription(description)
+                // This can be attached to the Status locally, but NOT 
transmitted to the client!
+                .withCause(cause)
+                .asRuntimeException();
+    }
+
     private void invokeRoute(GrpcEndpoint endpoint, Exchange exchange) throws 
Exception {
         if (endpoint.getConfiguration().isSynchronous()) {
             consumer.getProcessor().process(exchange);
diff --git 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAbstractStreamObserver.java
 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAbstractStreamObserver.java
index 2f501aa8ca42..adb2c6f8827c 100644
--- 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAbstractStreamObserver.java
+++ 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAbstractStreamObserver.java
@@ -40,4 +40,17 @@ public abstract class GrpcRequestAbstractStreamObserver 
implements StreamObserve
         this.responseObserver = responseObserver;
         this.headers = headers;
     }
+
+    /**
+     * Sends the failure of the exchange to the client as an error, like for 
unary calls.
+     *
+     * @return true if the exchange failed and the error was sent
+     */
+    protected boolean sendFailure(Exchange exchange) {
+        if (exchange.isFailed()) {
+            
responseObserver.onError(GrpcMethodHandler.toStatusException(endpoint, 
exchange.getException()));
+            return true;
+        }
+        return false;
+    }
 }
diff --git 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAggregationStreamObserver.java
 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAggregationStreamObserver.java
index d0aa4b2b1975..a56c4b5deaad 100644
--- 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAggregationStreamObserver.java
+++ 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestAggregationStreamObserver.java
@@ -62,6 +62,10 @@ public class GrpcRequestAggregationStreamObserver extends 
GrpcRequestAbstractStr
         try {
             latch.await();
 
+            if (sendFailure(exchange)) {
+                return;
+            }
+
             Object responseBody = exchange.getMessage().getBody();
             if (responseBody instanceof List) {
                 List<?> responseList = (List<?>) responseBody;
diff --git 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestPropagationStreamObserver.java
 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestPropagationStreamObserver.java
index 4764cceac0b7..a9e29cc114d9 100644
--- 
a/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestPropagationStreamObserver.java
+++ 
b/components/camel-grpc/src/main/java/org/apache/camel/component/grpc/server/GrpcRequestPropagationStreamObserver.java
@@ -35,8 +35,14 @@ public class GrpcRequestPropagationStreamObserver extends 
GrpcRequestAbstractStr
         super(endpoint, consumer, responseObserver, headers);
     }
 
+    // set when a failed exchange ended the call with an error: nothing more 
can be sent to the client
+    private volatile boolean failed;
+
     @Override
     public void onNext(Object request) {
+        if (failed) {
+            return;
+        }
         CountDownLatch latch = new CountDownLatch(1);
 
         exchange = endpoint.createExchange();
@@ -48,6 +54,11 @@ public class GrpcRequestPropagationStreamObserver extends 
GrpcRequestAbstractStr
         try {
             latch.await();
 
+            if (sendFailure(exchange)) {
+                failed = true;
+                return;
+            }
+
             Object responseBody = exchange.getMessage().getBody();
             if (responseBody instanceof List) {
                 List<?> responseList = (List<?>) responseBody;
@@ -58,7 +69,7 @@ public class GrpcRequestPropagationStreamObserver extends 
GrpcRequestAbstractStr
         } catch (InterruptedException e) {
             Thread.currentThread().interrupt();
             responseObserver.onError(e);
-
+            failed = true;
         }
     }
 
@@ -67,7 +78,9 @@ public class GrpcRequestPropagationStreamObserver extends 
GrpcRequestAbstractStr
         exchange = endpoint.createExchange();
         exchange.getIn().setHeaders(headers);
         consumer.onError(exchange, throwable);
-        responseObserver.onError(throwable);
+        if (!failed) {
+            responseObserver.onError(throwable);
+        }
     }
 
     @Override
@@ -75,6 +88,8 @@ public class GrpcRequestPropagationStreamObserver extends 
GrpcRequestAbstractStr
         exchange = endpoint.createExchange();
         exchange.getIn().setHeaders(headers);
         consumer.onCompleted(exchange);
-        responseObserver.onCompleted();
+        if (!failed) {
+            responseObserver.onCompleted();
+        }
     }
 }
diff --git 
a/components/camel-grpc/src/test/java/org/apache/camel/component/grpc/GrpcConsumerStreamingExceptionTest.java
 
b/components/camel-grpc/src/test/java/org/apache/camel/component/grpc/GrpcConsumerStreamingExceptionTest.java
new file mode 100644
index 000000000000..8244dbbb83a8
--- /dev/null
+++ 
b/components/camel-grpc/src/test/java/org/apache/camel/component/grpc/GrpcConsumerStreamingExceptionTest.java
@@ -0,0 +1,224 @@
+/*
+ * 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.grpc;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import io.grpc.ManagedChannel;
+import io.grpc.ManagedChannelBuilder;
+import io.grpc.Status;
+import io.grpc.StatusRuntimeException;
+import io.grpc.stub.StreamObserver;
+import org.apache.camel.CamelException;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * With the AGGREGATION and PROPAGATION consumer strategies a failed exchange 
must end the call with an error, as for
+ * unary calls, and not send the message body as a response.
+ */
+public class GrpcConsumerStreamingExceptionTest extends GrpcTestSupport {
+
+    private static final String ROUTE_EXCEPTION_MESSAGE = "GRPC Camel 
streaming exception message";
+    private static final String MUTED_EXCEPTION_MESSAGE = "Exchange processing 
failed";
+
+    private final Map<String, ManagedChannel> channels = new 
ConcurrentHashMap<>();
+
+    @AfterEach
+    public void stopGrpcChannels() {
+        channels.values().forEach(channel -> channel.shutdown().shutdownNow());
+        channels.clear();
+    }
+
+    @Test
+    public void testAggregationFailureIsSentAsError() throws Exception {
+        PongResponseStreamObserver responseObserver = 
callClientStreaming("grpc-aggregation");
+
+        assertTrue(responseObserver.pongResponses.isEmpty(), "no response must 
be sent for a failed exchange");
+        assertInternalError(responseObserver, ROUTE_EXCEPTION_MESSAGE);
+    }
+
+    @Test
+    public void testPropagationFailureIsSentAsError() throws Exception {
+        PongResponseStreamObserver responseObserver = 
callClientStreaming("grpc-propagation");
+
+        assertTrue(responseObserver.pongResponses.isEmpty(), "no response must 
be sent for a failed exchange");
+        assertInternalError(responseObserver, ROUTE_EXCEPTION_MESSAGE);
+    }
+
+    @Test
+    public void testPropagationFailureEndsTheCall() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:propagation");
+        mock.expectedMessageCount(1);
+        // a message routed after the failure would arrive later than the error
+        mock.setAssertPeriod(500);
+
+        PongResponseStreamObserver responseObserver = 
callBidiStreaming("grpc-propagation", 2);
+
+        assertTrue(responseObserver.pongResponses.isEmpty(), "no response must 
be sent for a failed exchange");
+        assertInternalError(responseObserver, ROUTE_EXCEPTION_MESSAGE);
+        mock.assertIsSatisfied();
+    }
+
+    @Test
+    public void testAggregationFailureIsMutedByDefault() throws Exception {
+        PongResponseStreamObserver responseObserver = 
callClientStreaming("grpc-aggregation-muted");
+
+        assertTrue(responseObserver.pongResponses.isEmpty(), "no response must 
be sent for a failed exchange");
+        assertInternalError(responseObserver, MUTED_EXCEPTION_MESSAGE);
+    }
+
+    @Test
+    public void testPropagationFailureIsMutedByDefault() throws Exception {
+        PongResponseStreamObserver responseObserver = 
callClientStreaming("grpc-propagation-muted");
+
+        assertTrue(responseObserver.pongResponses.isEmpty(), "no response must 
be sent for a failed exchange");
+        assertInternalError(responseObserver, MUTED_EXCEPTION_MESSAGE);
+    }
+
+    @Test
+    public void testPropagationHandledFailureKeepsTheStreamOpen() throws 
Exception {
+        PongResponseStreamObserver responseObserver = 
callBidiStreaming("grpc-propagation-handled", 2);
+
+        assertNull(responseObserver.error, "a handled exception must not end 
the call with an error");
+        assertEquals(1, responseObserver.completed.get());
+        List<Integer> pongIds = new ArrayList<>();
+        responseObserver.pongResponses.forEach(pong -> 
pongIds.add(pong.getPongId()));
+        assertEquals(List.of(1, 2), pongIds);
+    }
+
+    // client streaming call (one response) with a single message
+    private PongResponseStreamObserver callClientStreaming(String routeId) 
throws InterruptedException {
+        return call(routeId, 1, false);
+    }
+
+    // bidirectional streaming call, so that each message can be answered
+    private PongResponseStreamObserver callBidiStreaming(String routeId, int 
messages) throws InterruptedException {
+        return call(routeId, messages, true);
+    }
+
+    private PongResponseStreamObserver call(String routeId, int messages, 
boolean bidi) throws InterruptedException {
+        ManagedChannel channel = channels.computeIfAbsent(routeId,
+                id -> ManagedChannelBuilder.forAddress("localhost", 
getRoutePort(id)).usePlaintext().build());
+        PingPongGrpc.PingPongStub stub = PingPongGrpc.newStub(channel);
+        PongResponseStreamObserver responseObserver = new 
PongResponseStreamObserver();
+        StreamObserver<PingRequest> requestObserver
+                = bidi ? stub.pingAsyncAsync(responseObserver) : 
stub.pingAsyncSync(responseObserver);
+        for (int i = 1; i <= messages; i++) {
+            
requestObserver.onNext(PingRequest.newBuilder().setPingName("PING").setPingId(i).build());
+        }
+        requestObserver.onCompleted();
+        assertTrue(responseObserver.latch.await(5, TimeUnit.SECONDS));
+        return responseObserver;
+    }
+
+    private static void assertInternalError(PongResponseStreamObserver 
responseObserver, String description) {
+        StatusRuntimeException e = 
assertInstanceOf(StatusRuntimeException.class, responseObserver.error);
+        assertEquals(Status.Code.INTERNAL, e.getStatus().getCode());
+        assertEquals(description, e.getStatus().getDescription());
+        assertEquals(0, responseObserver.completed.get(), "a failed call must 
not also complete");
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("grpc://localhost:0/org.apache.camel.component.grpc.PingPong?synchronous=true"
+                     + "&consumerStrategy=AGGREGATION&muteException=false")
+                        .routeId("grpc-aggregation")
+                        .bean(new GrpcMessageBuilder(), 
"buildAggregatedPongResponse")
+                        .throwException(CamelException.class, 
ROUTE_EXCEPTION_MESSAGE);
+
+                
from("grpc://localhost:0/org.apache.camel.component.grpc.PingPong?synchronous=true"
+                     + "&consumerStrategy=PROPAGATION&muteException=false")
+                        .routeId("grpc-propagation")
+                        .bean(new GrpcMessageBuilder(), "buildPongResponse")
+                        .to("mock:propagation")
+                        .throwException(CamelException.class, 
ROUTE_EXCEPTION_MESSAGE);
+
+                
from("grpc://localhost:0/org.apache.camel.component.grpc.PingPong?synchronous=true"
+                     + "&consumerStrategy=AGGREGATION")
+                        .routeId("grpc-aggregation-muted")
+                        .bean(new GrpcMessageBuilder(), 
"buildAggregatedPongResponse")
+                        .throwException(CamelException.class, 
ROUTE_EXCEPTION_MESSAGE);
+
+                
from("grpc://localhost:0/org.apache.camel.component.grpc.PingPong?synchronous=true"
+                     + "&consumerStrategy=PROPAGATION")
+                        .routeId("grpc-propagation-muted")
+                        .bean(new GrpcMessageBuilder(), "buildPongResponse")
+                        .throwException(CamelException.class, 
ROUTE_EXCEPTION_MESSAGE);
+
+                
from("grpc://localhost:0/org.apache.camel.component.grpc.PingPong?synchronous=true"
+                     + "&consumerStrategy=PROPAGATION&muteException=true")
+                        .routeId("grpc-propagation-handled")
+                        .onException(CamelException.class).handled(true).end()
+                        .bean(new GrpcMessageBuilder(), "buildPongResponse")
+                        .throwException(CamelException.class, 
ROUTE_EXCEPTION_MESSAGE);
+            }
+        };
+    }
+
+    public static class GrpcMessageBuilder {
+        public PongResponse buildAggregatedPongResponse(List<PingRequest> 
pingRequests) {
+            return buildPongResponse(pingRequests.get(0));
+        }
+
+        public PongResponse buildPongResponse(PingRequest pingRequest) {
+            return 
PongResponse.newBuilder().setPongName(pingRequest.getPingName() + "PONG")
+                    .setPongId(pingRequest.getPingId()).build();
+        }
+    }
+
+    static class PongResponseStreamObserver implements 
StreamObserver<PongResponse> {
+        private final CountDownLatch latch = new CountDownLatch(1);
+        private final List<PongResponse> pongResponses = new 
CopyOnWriteArrayList<>();
+        private final AtomicInteger completed = new AtomicInteger();
+        private volatile Throwable error;
+
+        @Override
+        public void onNext(PongResponse value) {
+            pongResponses.add(value);
+        }
+
+        @Override
+        public void onError(Throwable t) {
+            error = t;
+            latch.countDown();
+        }
+
+        @Override
+        public void onCompleted() {
+            completed.incrementAndGet();
+            latch.countDown();
+        }
+    }
+}
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 66ea24ac74a5..37d5b442c237 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
@@ -2266,6 +2266,15 @@ A route that relies on the exception message reaching 
the client must opt back i
 grpc://localhost:8080/org.example.MyService?muteException=false
 ----
 
+==== Failed exchanges with the AGGREGATION and PROPAGATION consumer strategies
+
+With the `AGGREGATION` and `PROPAGATION` consumer strategies (client streaming 
and bidirectional streaming calls), a
+failed exchange now ends the call with an `INTERNAL` error, like unary and 
server streaming calls already did
+(the description follows the `muteException` option). Before, the message body 
of the failed exchange was sent to
+the client as a regular response. With `PROPAGATION` (the default strategy) 
the first failed exchange ends the whole
+call, so the messages that the client sends afterwards are no longer routed; a 
route that must keep the stream open
+can handle the exception (`onException(...).handled(true)`) and answer that 
message with a body.
+
 === camel-keycloak
 
 When `validateIssuer` is enabled and the token is checked by introspection, an 
introspection response

Reply via email to