From 9d14356cea87c2cc8cd44172ea3c2dab3074832f Mon Sep 17 00:00:00 2001 From: AgraVator Date: Wed, 22 Jul 2026 20:01:28 +0530 Subject: [PATCH 1/2] outlier-detection: exclude client and hedging cancellations from call counters --- .../main/java/io/grpc/ClientStreamTracer.java | 9 ++ .../grpc/internal/AbstractClientStream.java | 1 + .../ForwardingClientStreamTracer.java | 5 + .../io/grpc/internal/StatsTraceContext.java | 13 ++ .../grpc/internal/StatsTraceContextTest.java | 49 ++++++++ .../util/ForwardingClientStreamTracer.java | 5 + .../util/OutlierDetectionLoadBalancer.java | 23 +++- .../OutlierDetectionLoadBalancerTest.java | 114 ++++++++++++++++++ 8 files changed, 217 insertions(+), 2 deletions(-) create mode 100644 core/src/test/java/io/grpc/internal/StatsTraceContextTest.java diff --git a/api/src/main/java/io/grpc/ClientStreamTracer.java b/api/src/main/java/io/grpc/ClientStreamTracer.java index 8e11e781e7c..71f4145eb82 100644 --- a/api/src/main/java/io/grpc/ClientStreamTracer.java +++ b/api/src/main/java/io/grpc/ClientStreamTracer.java @@ -99,6 +99,15 @@ public void inboundTrailers(Metadata trailers) { public void addOptionalLabel(String key, String value) { } + /** + * The stream was cancelled from the client side before a normal response was received. + * + * @param status the cancellation status + * @since 1.70.0 + */ + public void cancelled(Status status) { + } + /** * Factory class for {@link ClientStreamTracer}. */ diff --git a/core/src/main/java/io/grpc/internal/AbstractClientStream.java b/core/src/main/java/io/grpc/internal/AbstractClientStream.java index bce1820b482..39e457f815f 100644 --- a/core/src/main/java/io/grpc/internal/AbstractClientStream.java +++ b/core/src/main/java/io/grpc/internal/AbstractClientStream.java @@ -198,6 +198,7 @@ public final void halfClose() { public final void cancel(Status reason) { Preconditions.checkArgument(!reason.isOk(), "Should not cancel with OK status"); cancelled = true; + transportState().getStatsTraceContext().clientCancelled(reason); abstractClientStreamSink().cancel(reason); } diff --git a/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java b/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java index e7679ea14cc..6dc18fdf627 100644 --- a/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java +++ b/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java @@ -64,6 +64,11 @@ public void addOptionalLabel(String key, String value) { delegate().addOptionalLabel(key, value); } + @Override + public void cancelled(Status status) { + delegate().cancelled(status); + } + @Override public void streamClosed(Status status) { delegate().streamClosed(status); diff --git a/core/src/main/java/io/grpc/internal/StatsTraceContext.java b/core/src/main/java/io/grpc/internal/StatsTraceContext.java index 007aefc0fb8..2827f6a9766 100644 --- a/core/src/main/java/io/grpc/internal/StatsTraceContext.java +++ b/core/src/main/java/io/grpc/internal/StatsTraceContext.java @@ -167,6 +167,19 @@ public void serverCallMethodResolved(MethodDescriptor method) { } } + /** + * See {@link ClientStreamTracer#cancelled}. For client-side only. + * + *

Called from abstract stream implementations. + */ + public void clientCancelled(Status status) { + for (StreamTracer tracer : tracers) { + if (tracer instanceof ClientStreamTracer) { + ((ClientStreamTracer) tracer).cancelled(status); + } + } + } + /** * See {@link StreamTracer#streamClosed}. This may be called multiple times, and only the first * value will be taken. diff --git a/core/src/test/java/io/grpc/internal/StatsTraceContextTest.java b/core/src/test/java/io/grpc/internal/StatsTraceContextTest.java new file mode 100644 index 00000000000..c00efd9e8de --- /dev/null +++ b/core/src/test/java/io/grpc/internal/StatsTraceContextTest.java @@ -0,0 +1,49 @@ +/* + * Copyright 2026 The gRPC Authors + * + * Licensed 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 io.grpc.internal; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; + +import io.grpc.ClientStreamTracer; +import io.grpc.ServerStreamTracer; +import io.grpc.Status; +import io.grpc.StreamTracer; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** Unit tests for {@link StatsTraceContext}. */ +@RunWith(JUnit4.class) +public class StatsTraceContextTest { + + @Test + public void clientCancelled_notifiesClientStreamTracers() { + ClientStreamTracer clientTracer = mock(ClientStreamTracer.class); + ServerStreamTracer serverTracer = mock(ServerStreamTracer.class); + + StatsTraceContext statsTraceCtx = new StatsTraceContext( + new StreamTracer[] {clientTracer, serverTracer}); + + Status cancelledStatus = Status.CANCELLED.withDescription("Client cancelled"); + statsTraceCtx.clientCancelled(cancelledStatus); + + verify(clientTracer).cancelled(cancelledStatus); + verifyNoInteractions(serverTracer); + } +} diff --git a/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java b/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java index 9c9998571e5..1eda7c4bdd2 100644 --- a/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java +++ b/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java @@ -63,6 +63,11 @@ public void addOptionalLabel(String key, String value) { delegate().addOptionalLabel(key, value); } + @Override + public void cancelled(Status status) { + delegate().cancelled(status); + } + @Override public void streamClosed(Status status) { delegate().streamClosed(status); diff --git a/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java b/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java index dc61441bccd..b7fd50ae6cd 100644 --- a/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java +++ b/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java @@ -477,22 +477,41 @@ public ClientStreamTracer newClientStreamTracer(StreamInfo info, Metadata header if (delegateFactory != null) { ClientStreamTracer delegateTracer = delegateFactory.newClientStreamTracer(info, headers); return new ForwardingClientStreamTracer() { + private volatile boolean cancelled; + @Override protected ClientStreamTracer delegate() { return delegateTracer; } + @Override + public void cancelled(Status status) { + cancelled = true; + delegate().cancelled(status); + } + @Override public void streamClosed(Status status) { - tracker.incrementCallCount(status.isOk()); + if (!cancelled) { + tracker.incrementCallCount(status.isOk()); + } delegate().streamClosed(status); } }; } else { return new ClientStreamTracer() { + private volatile boolean cancelled; + + @Override + public void cancelled(Status status) { + cancelled = true; + } + @Override public void streamClosed(Status status) { - tracker.incrementCallCount(status.isOk()); + if (!cancelled) { + tracker.incrementCallCount(status.isOk()); + } } }; } diff --git a/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java b/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java index 39f5b5fb7d6..c359a316618 100644 --- a/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java +++ b/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java @@ -451,6 +451,29 @@ public void delegatePickTracerFactoryPreserved() { verify(mockStreamTracer).inboundHeaders(); } + @Test + public void delegatePick_cancelledForwarded() { + OutlierDetectionLoadBalancerConfig config = new OutlierDetectionLoadBalancerConfig.Builder() + .setSuccessRateEjection(new SuccessRateEjection.Builder().build()) + .setChildConfig(newChildConfig(fakeLbProvider, null)).build(); + + loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers.get(0))); + + final Subchannel readySubchannel = subchannels.values().iterator().next(); + deliverSubchannelState(readySubchannel, ConnectivityStateInfo.forNonError(READY)); + + verify(mockHelper, times(2)).updateBalancingState(stateCaptor.capture(), + pickerCaptor.capture()); + + SubchannelPicker picker = pickerCaptor.getAllValues().get(1); + PickResult pickResult = picker.pickSubchannel(mock(PickSubchannelArgs.class)); + + ClientStreamTracer clientStreamTracer = pickResult.getStreamTracerFactory() + .newClientStreamTracer(ClientStreamTracer.StreamInfo.newBuilder().build(), new Metadata()); + clientStreamTracer.cancelled(Status.CANCELLED); + verify(mockStreamTracer).cancelled(Status.CANCELLED); + } + /** * Assure the tracer works even when the underlying LB does not have a tracer to delegate to. */ @@ -531,6 +554,51 @@ public void successRateOneOutlier() { assertEjectedSubchannels(ImmutableSet.of(ImmutableSet.copyOf(servers.get(0).getAddresses()))); } + /** + * Client-cancelled streams (e.g. non-winning hedged attempts) do not count as failures. + */ + @Test + public void successRate_clientCancelled_notEjected() { + OutlierDetectionLoadBalancerConfig config = new OutlierDetectionLoadBalancerConfig.Builder() + .setMaxEjectionPercent(50) + .setSuccessRateEjection( + new SuccessRateEjection.Builder() + .setMinimumHosts(3) + .setRequestVolume(10).build()) + .setChildConfig(newChildConfig(roundRobinLbProvider, null)).build(); + + loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); + + deliverSubchannelState(subchannel1, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel2, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel3, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel4, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel5, ConnectivityStateInfo.forNonError(READY)); + + verify(mockHelper, times(7)).updateBalancingState(stateCaptor.capture(), + pickerCaptor.capture()); + SubchannelPicker picker = pickerCaptor.getAllValues() + .get(pickerCaptor.getAllValues().size() - 1); + + for (int i = 0; i < 100; i++) { + PickResult pickResult = picker.pickSubchannel(mock(PickSubchannelArgs.class)); + ClientStreamTracer clientStreamTracer = pickResult.getStreamTracerFactory() + .newClientStreamTracer(null, null); + Subchannel subchannel = (Subchannel) pickResult.getSubchannel().getInternalSubchannel(); + if (subchannel == subchannel1) { + clientStreamTracer.cancelled(Status.CANCELLED); + clientStreamTracer.streamClosed(Status.CANCELLED); + } else { + clientStreamTracer.streamClosed(Status.OK); + } + } + + forwardTime(config); + + // subchannel1 was cancelled client-side and should not be ejected as an outlier. + assertEjectedSubchannels(ImmutableSet.of()); + } + /** * The success rate algorithm ejects the outlier, but then the config changes so that similar * behavior no longer gets ejected. @@ -781,6 +849,52 @@ public void failurePercentageNoOutliers() { assertEjectedSubchannels(ImmutableSet.of()); } + /** + * Client-cancelled streams (e.g. non-winning hedged attempts) do not count as failures for + * failure percentage algorithm. + */ + @Test + public void failurePercentage_clientCancelled_notEjected() { + OutlierDetectionLoadBalancerConfig config = new OutlierDetectionLoadBalancerConfig.Builder() + .setMaxEjectionPercent(50) + .setFailurePercentageEjection( + new FailurePercentageEjection.Builder() + .setMinimumHosts(3) + .setRequestVolume(10).build()) + .setChildConfig(newChildConfig(roundRobinLbProvider, null)).build(); + + loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); + + deliverSubchannelState(subchannel1, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel2, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel3, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel4, ConnectivityStateInfo.forNonError(READY)); + deliverSubchannelState(subchannel5, ConnectivityStateInfo.forNonError(READY)); + + verify(mockHelper, times(7)).updateBalancingState(stateCaptor.capture(), + pickerCaptor.capture()); + SubchannelPicker picker = pickerCaptor.getAllValues() + .get(pickerCaptor.getAllValues().size() - 1); + + for (int i = 0; i < 100; i++) { + PickResult pickResult = picker.pickSubchannel(mock(PickSubchannelArgs.class)); + ClientStreamTracer clientStreamTracer = pickResult.getStreamTracerFactory() + .newClientStreamTracer(null, null); + Subchannel subchannel = (Subchannel) pickResult.getSubchannel().getInternalSubchannel(); + if (subchannel == subchannel1) { + clientStreamTracer.cancelled(Status.CANCELLED); + clientStreamTracer.streamClosed(Status.CANCELLED); + } else { + clientStreamTracer.streamClosed(Status.OK); + } + } + + forwardTime(config); + + // subchannel1 was cancelled client-side and should not be ejected as an outlier. + assertEjectedSubchannels(ImmutableSet.of()); + } + /** * The success rate algorithm ejects the outlier. */ From 7abb48523467ee97e53f05c3a94218bb9b6013c1 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Wed, 22 Jul 2026 20:08:59 +0530 Subject: [PATCH 2/2] Update @since tag to 1.84.0 and add tracer cancellation in FailingClientStream --- .../main/java/io/grpc/ClientStreamTracer.java | 2 +- .../io/grpc/internal/FailingClientStream.java | 7 +++++++ .../internal/AbstractClientStreamTest.java | 19 +++++++++++++++++++ .../internal/FailingClientStreamTest.java | 10 ++++++++++ 4 files changed, 37 insertions(+), 1 deletion(-) diff --git a/api/src/main/java/io/grpc/ClientStreamTracer.java b/api/src/main/java/io/grpc/ClientStreamTracer.java index 71f4145eb82..537457d2a9e 100644 --- a/api/src/main/java/io/grpc/ClientStreamTracer.java +++ b/api/src/main/java/io/grpc/ClientStreamTracer.java @@ -103,7 +103,7 @@ public void addOptionalLabel(String key, String value) { * The stream was cancelled from the client side before a normal response was received. * * @param status the cancellation status - * @since 1.70.0 + * @since 1.84.0 */ public void cancelled(Status status) { } diff --git a/core/src/main/java/io/grpc/internal/FailingClientStream.java b/core/src/main/java/io/grpc/internal/FailingClientStream.java index 6388ef8b6ee..a81be05ba43 100644 --- a/core/src/main/java/io/grpc/internal/FailingClientStream.java +++ b/core/src/main/java/io/grpc/internal/FailingClientStream.java @@ -61,6 +61,13 @@ public void start(ClientStreamListener listener) { listener.closed(error, rpcProgress, new Metadata()); } + @Override + public void cancel(Status reason) { + for (ClientStreamTracer tracer : tracers) { + tracer.cancelled(reason); + } + } + @VisibleForTesting Status getError() { return error; diff --git a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java index 8f14b74035c..c985605fda3 100644 --- a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java @@ -39,6 +39,7 @@ import io.grpc.Attributes; import io.grpc.CallOptions; +import io.grpc.ClientStreamTracer; import io.grpc.Codec; import io.grpc.Deadline; import io.grpc.Grpc; @@ -155,6 +156,24 @@ public void cancel(Status errorStatus) { verify(mockListener).closed(any(Status.class), same(PROCESSED), any(Metadata.class)); } + @Test + public void cancel_notifiesStatsTraceContext() { + ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); + StatsTraceContext customStatsTraceCtx = new StatsTraceContext(new StreamTracer[] {mockTracer}); + final BaseTransportState state = new BaseTransportState(customStatsTraceCtx, transportTracer); + AbstractClientStream stream = new BaseAbstractClientStream(allocator, state, new BaseSink() { + @Override + public void cancel(Status errorStatus) { + } + }, customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + Status cancelStatus = Status.CANCELLED.withDescription("Cancelled by test"); + stream.cancel(cancelStatus); + + verify(mockTracer).cancelled(cancelStatus); + } + @Test public void startFailsOnNullListener() { AbstractClientStream stream = diff --git a/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java b/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java index c07812577d5..828388c4fbd 100644 --- a/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java @@ -57,4 +57,14 @@ public void droppedRpcProgressPopulatedToListener() { stream.start(listener); verify(listener).closed(eq(status), eq(RpcProgress.DROPPED), any(Metadata.class)); } + + @Test + public void cancel_notifiesTracers() { + ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); + ClientStream stream = new FailingClientStream( + Status.UNAVAILABLE, RpcProgress.PROCESSED, new ClientStreamTracer[] {mockTracer}); + Status cancelStatus = Status.CANCELLED.withDescription("Cancelled by test"); + stream.cancel(cancelStatus); + verify(mockTracer).cancelled(cancelStatus); + } }