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