From 9d14356cea87c2cc8cd44172ea3c2dab3074832f Mon Sep 17 00:00:00 2001 From: AgraVator Date: Wed, 22 Jul 2026 20:01:28 +0530 Subject: [PATCH 01/15] 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 02/15] 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); + } } From f9cc87f65dd645a6713f5fe90682b277c2024c47 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Thu, 23 Jul 2026 15:20:39 +0530 Subject: [PATCH 03/15] core, inprocess, binder: synchronize client cancellation tracer notifications with transport stream closure --- .../src/main/java/io/grpc/binder/internal/Inbound.java | 3 +++ .../java/io/grpc/internal/AbstractClientStream.java | 4 +++- .../java/io/grpc/internal/FailingClientStream.java | 7 ------- .../io/grpc/internal/AbstractClientStreamTest.java | 1 + .../java/io/grpc/internal/FailingClientStreamTest.java | 10 ---------- .../java/io/grpc/inprocess/InProcessTransport.java | 1 + 6 files changed, 8 insertions(+), 18 deletions(-) diff --git a/binder/src/main/java/io/grpc/binder/internal/Inbound.java b/binder/src/main/java/io/grpc/binder/internal/Inbound.java index 83fc8273d6f..7eea5ee78a1 100644 --- a/binder/src/main/java/io/grpc/binder/internal/Inbound.java +++ b/binder/src/main/java/io/grpc/binder/internal/Inbound.java @@ -268,6 +268,9 @@ private final void deliverInternal() { @GuardedBy("this") final void closeOnCancel(Status status) { + if (!isClosed() && statsTraceContext != null) { + statsTraceContext.clientCancelled(status); + } closeAbnormal(Status.CANCELLED, status, false); } diff --git a/core/src/main/java/io/grpc/internal/AbstractClientStream.java b/core/src/main/java/io/grpc/internal/AbstractClientStream.java index 39e457f815f..e6cd76c99b9 100644 --- a/core/src/main/java/io/grpc/internal/AbstractClientStream.java +++ b/core/src/main/java/io/grpc/internal/AbstractClientStream.java @@ -198,7 +198,6 @@ 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); } @@ -458,6 +457,9 @@ private void closeListener( Status status, RpcProgress rpcProgress, Metadata trailers) { if (!listenerClosed) { listenerClosed = true; + if (status.getCode() == Status.Code.CANCELLED) { + statsTraceCtx.clientCancelled(status); + } statsTraceCtx.streamClosed(status); if (getTransportTracer() != null) { getTransportTracer().reportStreamClosed(status.isOk()); diff --git a/core/src/main/java/io/grpc/internal/FailingClientStream.java b/core/src/main/java/io/grpc/internal/FailingClientStream.java index a81be05ba43..6388ef8b6ee 100644 --- a/core/src/main/java/io/grpc/internal/FailingClientStream.java +++ b/core/src/main/java/io/grpc/internal/FailingClientStream.java @@ -61,13 +61,6 @@ 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 c985605fda3..a2ebf28ff92 100644 --- a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java @@ -164,6 +164,7 @@ public void cancel_notifiesStatsTraceContext() { AbstractClientStream stream = new BaseAbstractClientStream(allocator, state, new BaseSink() { @Override public void cancel(Status errorStatus) { + state.transportReportStatus(errorStatus, true, new Metadata()); } }, customStatsTraceCtx, transportTracer); stream.start(mockListener); diff --git a/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java b/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java index 828388c4fbd..c07812577d5 100644 --- a/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/FailingClientStreamTest.java @@ -57,14 +57,4 @@ 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); - } } diff --git a/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java b/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java index a92f10fd5c5..6c11229fa3d 100644 --- a/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java +++ b/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java @@ -846,6 +846,7 @@ public void cancel(Status reason) { if (!internalCancel(serverStatus, serverStatus)) { return; } + statsTraceCtx.clientCancelled(reason); serverStream.clientCancelled(reason); streamClosed(); } From 120e35ccad708a014b47f653f8981eeef0886c8f Mon Sep 17 00:00:00 2001 From: AgraVator Date: Fri, 24 Jul 2026 11:36:17 +0530 Subject: [PATCH 04/15] test: add unit tests for race condition and transport tracer cancellation --- .../internal/AbstractClientStreamTest.java | 22 +++++++++ .../inprocess/InProcessTransportTest.java | 48 +++++++++++++++++++ 2 files changed, 70 insertions(+) diff --git a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java index a2ebf28ff92..346b6fe6e15 100644 --- a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java @@ -175,6 +175,28 @@ public void cancel(Status errorStatus) { verify(mockTracer).cancelled(cancelStatus); } + @Test + public void transportReportStatus_okFirst_lateCancellationDoesNotNotifyTracerCancelled() { + 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) { + state.transportReportStatus(errorStatus, true, new Metadata()); + } + }, customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + // Report Status.OK first + state.transportReportStatus(Status.OK, false, new Metadata()); + + // Subsequent late cancellation + stream.cancel(Status.CANCELLED.withDescription("Late cancel")); + + verify(mockTracer, never()).cancelled(any(Status.class)); + } + @Test public void startFailsOnNullListener() { AbstractClientStream stream = diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index d2220e05114..d1017cdb0eb 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -236,6 +236,54 @@ public void basicStreamInProcess() throws Exception { serverStream.close(status, new Metadata()); } + @Test + public void clientStream_cancel_notifiesTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener = + serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), CallOptions.DEFAULT, + new io.grpc.ClientStreamTracer[] {mockTracer}); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + + Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled"); + clientStream.cancel(cancelStatus); + + org.mockito.Mockito.verify(mockTracer).cancelled(cancelStatus); + } + + @Test + public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener = + serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), CallOptions.DEFAULT, + new io.grpc.ClientStreamTracer[] {mockTracer}); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + StreamCreation serverStreamCreation = + serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + ServerStream serverStream = serverStreamCreation.stream; + + serverStream.close(Status.OK, new Metadata()); + clientStream.cancel(Status.CANCELLED.withDescription("Late cancellation")); + + org.mockito.Mockito.verify(mockTracer, org.mockito.Mockito.never()) + .cancelled(org.mockito.Mockito.any(Status.class)); + } + private void assertAssumedMessageSize( TestStreamTracer streamTracerSender, TestStreamTracer streamTracerReceiver) { if (isEnabledSupportTracingMessageSizes()) { From f9dd0fb865052ba855221ad455b6b70e0bfdc855 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Fri, 24 Jul 2026 11:39:48 +0530 Subject: [PATCH 05/15] test: add inprocess client stream cancellation unit tests --- .../test/java/io/grpc/inprocess/InProcessTransportTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index d1017cdb0eb..0cecb72b9de 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -238,7 +238,6 @@ public void basicStreamInProcess() throws Exception { @Test public void clientStream_cancel_notifiesTracerCancelled() throws Exception { - server = newServer(Arrays.asList(serverStreamTracerFactory)); client = newClientTransport(server); startTransport(client, mockClientTransportListener); MockServerTransportListener serverTransportListener = @@ -251,6 +250,8 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, new io.grpc.ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); + StreamCreation serverStreamCreation = + serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled"); clientStream.cancel(cancelStatus); @@ -260,7 +261,6 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, @Test public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() throws Exception { - server = newServer(Arrays.asList(serverStreamTracerFactory)); client = newClientTransport(server); startTransport(client, mockClientTransportListener); MockServerTransportListener serverTransportListener = From fe581ae3e12e1b69f76fefb29236fa3f949cc347 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Fri, 24 Jul 2026 13:15:55 +0530 Subject: [PATCH 06/15] Revert "test: add inprocess client stream cancellation unit tests" This reverts commit f9dd0fb865052ba855221ad455b6b70e0bfdc855. --- .../test/java/io/grpc/inprocess/InProcessTransportTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index 0cecb72b9de..d1017cdb0eb 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -238,6 +238,7 @@ public void basicStreamInProcess() throws Exception { @Test public void clientStream_cancel_notifiesTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); client = newClientTransport(server); startTransport(client, mockClientTransportListener); MockServerTransportListener serverTransportListener = @@ -250,8 +251,6 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, new io.grpc.ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); - StreamCreation serverStreamCreation = - serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled"); clientStream.cancel(cancelStatus); @@ -261,6 +260,7 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, @Test public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); client = newClientTransport(server); startTransport(client, mockClientTransportListener); MockServerTransportListener serverTransportListener = From 5e558a5bf2606cb2c8c2beceb4c8369ef2e1c189 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Fri, 24 Jul 2026 14:09:34 +0530 Subject: [PATCH 07/15] test: fix checkstyle line length violations in InProcessTransportTest --- .../inprocess/InProcessTransportTest.java | 21 ++++++++++++------- 1 file changed, 13 insertions(+), 8 deletions(-) diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index d1017cdb0eb..dcb5d74b3ff 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -19,9 +19,14 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; import io.grpc.CallOptions; import io.grpc.ClientCall; +import io.grpc.ClientStreamTracer; import io.grpc.ManagedChannel; import io.grpc.Metadata; import io.grpc.MethodDescriptor; @@ -245,21 +250,22 @@ public void clientStream_cancel_notifiesTracerCancelled() throws Exception { serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); serverTransport = serverTransportListener.transport; - io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); + ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); ClientStream clientStream = client.newStream( methodDescriptor, new Metadata(), CallOptions.DEFAULT, - new io.grpc.ClientStreamTracer[] {mockTracer}); + new ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled"); clientStream.cancel(cancelStatus); - org.mockito.Mockito.verify(mockTracer).cancelled(cancelStatus); + verify(mockTracer).cancelled(cancelStatus); } @Test - public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() throws Exception { + public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() + throws Exception { server = newServer(Arrays.asList(serverStreamTracerFactory)); client = newClientTransport(server); startTransport(client, mockClientTransportListener); @@ -267,10 +273,10 @@ public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() t serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); serverTransport = serverTransportListener.transport; - io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); + ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); ClientStream clientStream = client.newStream( methodDescriptor, new Metadata(), CallOptions.DEFAULT, - new io.grpc.ClientStreamTracer[] {mockTracer}); + new ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); StreamCreation serverStreamCreation = @@ -280,8 +286,7 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, serverStream.close(Status.OK, new Metadata()); clientStream.cancel(Status.CANCELLED.withDescription("Late cancellation")); - org.mockito.Mockito.verify(mockTracer, org.mockito.Mockito.never()) - .cancelled(org.mockito.Mockito.any(Status.class)); + verify(mockTracer, never()).cancelled(any(Status.class)); } private void assertAssumedMessageSize( From 634be8136fc867dc2c3abdd9f0c45bda36645c17 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Fri, 24 Jul 2026 14:25:28 +0530 Subject: [PATCH 08/15] Revert "test: fix checkstyle line length violations in InProcessTransportTest" This reverts commit 5e558a5bf2606cb2c8c2beceb4c8369ef2e1c189. --- .../inprocess/InProcessTransportTest.java | 21 +++++++------------ 1 file changed, 8 insertions(+), 13 deletions(-) diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index dcb5d74b3ff..d1017cdb0eb 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -19,14 +19,9 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; -import static org.mockito.Mockito.any; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.never; -import static org.mockito.Mockito.verify; import io.grpc.CallOptions; import io.grpc.ClientCall; -import io.grpc.ClientStreamTracer; import io.grpc.ManagedChannel; import io.grpc.Metadata; import io.grpc.MethodDescriptor; @@ -250,22 +245,21 @@ public void clientStream_cancel_notifiesTracerCancelled() throws Exception { serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); serverTransport = serverTransportListener.transport; - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); + io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); ClientStream clientStream = client.newStream( methodDescriptor, new Metadata(), CallOptions.DEFAULT, - new ClientStreamTracer[] {mockTracer}); + new io.grpc.ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled"); clientStream.cancel(cancelStatus); - verify(mockTracer).cancelled(cancelStatus); + org.mockito.Mockito.verify(mockTracer).cancelled(cancelStatus); } @Test - public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() - throws Exception { + public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() throws Exception { server = newServer(Arrays.asList(serverStreamTracerFactory)); client = newClientTransport(server); startTransport(client, mockClientTransportListener); @@ -273,10 +267,10 @@ public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); serverTransport = serverTransportListener.transport; - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); + io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); ClientStream clientStream = client.newStream( methodDescriptor, new Metadata(), CallOptions.DEFAULT, - new ClientStreamTracer[] {mockTracer}); + new io.grpc.ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); StreamCreation serverStreamCreation = @@ -286,7 +280,8 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, serverStream.close(Status.OK, new Metadata()); clientStream.cancel(Status.CANCELLED.withDescription("Late cancellation")); - verify(mockTracer, never()).cancelled(any(Status.class)); + org.mockito.Mockito.verify(mockTracer, org.mockito.Mockito.never()) + .cancelled(org.mockito.Mockito.any(Status.class)); } private void assertAssumedMessageSize( From 20682d91fd777c8cfd4b9dc01d82f74c10f5d759 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Mon, 27 Jul 2026 12:55:07 +0530 Subject: [PATCH 09/15] test: fix checkstyle violations and missing server.start in InProcessTransportTest --- .../io/grpc/inprocess/InProcessTransportTest.java | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index d1017cdb0eb..1f977eaecb5 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -22,6 +22,7 @@ import io.grpc.CallOptions; import io.grpc.ClientCall; +import io.grpc.ClientStreamTracer; import io.grpc.ManagedChannel; import io.grpc.Metadata; import io.grpc.MethodDescriptor; @@ -239,16 +240,17 @@ public void basicStreamInProcess() throws Exception { @Test public void clientStream_cancel_notifiesTracerCancelled() throws Exception { server = newServer(Arrays.asList(serverStreamTracerFactory)); + server.start(serverListener); client = newClientTransport(server); startTransport(client, mockClientTransportListener); MockServerTransportListener serverTransportListener = serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); serverTransport = serverTransportListener.transport; - io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); + ClientStreamTracer mockTracer = org.mockito.Mockito.mock(ClientStreamTracer.class); ClientStream clientStream = client.newStream( methodDescriptor, new Metadata(), CallOptions.DEFAULT, - new io.grpc.ClientStreamTracer[] {mockTracer}); + new ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); @@ -261,16 +263,17 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, @Test public void clientStream_cancelAfterServerClose_doesNotNotifyTracerCancelled() throws Exception { server = newServer(Arrays.asList(serverStreamTracerFactory)); + server.start(serverListener); client = newClientTransport(server); startTransport(client, mockClientTransportListener); MockServerTransportListener serverTransportListener = serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); serverTransport = serverTransportListener.transport; - io.grpc.ClientStreamTracer mockTracer = org.mockito.Mockito.mock(io.grpc.ClientStreamTracer.class); + ClientStreamTracer mockTracer = org.mockito.Mockito.mock(ClientStreamTracer.class); ClientStream clientStream = client.newStream( methodDescriptor, new Metadata(), CallOptions.DEFAULT, - new io.grpc.ClientStreamTracer[] {mockTracer}); + new ClientStreamTracer[] {mockTracer}); ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); clientStream.start(clientStreamListener); StreamCreation serverStreamCreation = From bf90eb5a69b901934104e8d4972daf4c63c0d581 Mon Sep 17 00:00:00 2001 From: agrawalabhi Date: Mon, 3 Aug 2026 11:08:49 +0000 Subject: [PATCH 10/15] util: exclude client cancellations and deadline exceeded from outlier detection call counter --- .../grpc/internal/AbstractClientStream.java | 11 +- .../util/OutlierDetectionLoadBalancer.java | 8 +- .../OutlierDetectionLoadBalancerTest.java | 165 ++++++++++++------ 3 files changed, 124 insertions(+), 60 deletions(-) diff --git a/core/src/main/java/io/grpc/internal/AbstractClientStream.java b/core/src/main/java/io/grpc/internal/AbstractClientStream.java index e6cd76c99b9..b3995b5380c 100644 --- a/core/src/main/java/io/grpc/internal/AbstractClientStream.java +++ b/core/src/main/java/io/grpc/internal/AbstractClientStream.java @@ -197,7 +197,11 @@ public final void halfClose() { @Override public final void cancel(Status reason) { Preconditions.checkArgument(!reason.isOk(), "Should not cancel with OK status"); + if (cancelled || transportState().isListenerClosed()) { + return; + } cancelled = true; + transportState().getStatsTraceContext().clientCancelled(reason); abstractClientStreamSink().cancel(reason); } @@ -251,6 +255,10 @@ protected TransportState( } } + protected final boolean isListenerClosed() { + return listenerClosed; + } + private void setFullStreamDecompression(boolean fullStreamDecompression) { this.fullStreamDecompression = fullStreamDecompression; } @@ -457,9 +465,6 @@ private void closeListener( Status status, RpcProgress rpcProgress, Metadata trailers) { if (!listenerClosed) { listenerClosed = true; - if (status.getCode() == Status.Code.CANCELLED) { - statsTraceCtx.clientCancelled(status); - } statsTraceCtx.streamClosed(status); if (getTransportTracer() != null) { getTransportTracer().reportStreamClosed(status.isOk()); diff --git a/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java b/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java index b7fd50ae6cd..966affc6cff 100644 --- a/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java +++ b/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java @@ -493,7 +493,7 @@ public void cancelled(Status status) { @Override public void streamClosed(Status status) { if (!cancelled) { - tracker.incrementCallCount(status.isOk()); + tracker.incrementCallCount(status); } delegate().streamClosed(status); } @@ -510,7 +510,7 @@ public void cancelled(Status status) { @Override public void streamClosed(Status status) { if (!cancelled) { - tracker.incrementCallCount(status.isOk()); + tracker.incrementCallCount(status); } } }; @@ -571,13 +571,13 @@ Set getSubchannels() { return ImmutableSet.copyOf(subchannels); } - void incrementCallCount(boolean success) { + void incrementCallCount(Status status) { // If neither algorithm is configured, no point in incrementing counters. if (config.successRateEjection == null && config.failurePercentageEjection == null) { return; } - if (success) { + if (status.isOk()) { activeCallCounter.successCount.getAndIncrement(); } else { activeCallCounter.failureCount.getAndIncrement(); diff --git a/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java b/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java index c359a316618..0e008d3dfce 100644 --- a/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java +++ b/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java @@ -545,7 +545,7 @@ public void successRateOneOutlier() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -555,10 +555,10 @@ public void successRateOneOutlier() { } /** - * Client-cancelled streams (e.g. non-winning hedged attempts) do not count as failures. + * The success rate algorithm ignores CANCELLED status calls. */ @Test - public void successRate_clientCancelled_notEjected() { + public void successRateOneOutlier_cancelledIgnored() { OutlierDetectionLoadBalancerConfig config = new OutlierDetectionLoadBalancerConfig.Builder() .setMaxEjectionPercent(50) .setSuccessRateEjection( @@ -569,33 +569,38 @@ public void successRate_clientCancelled_notEjected() { 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)); + // subchannel1 returns CANCELLED. + generateLoad(ImmutableMap.of(subchannel1, Status.CANCELLED), 7); - verify(mockHelper, times(7)).updateBalancingState(stateCaptor.capture(), - pickerCaptor.capture()); - SubchannelPicker picker = pickerCaptor.getAllValues() - .get(pickerCaptor.getAllValues().size() - 1); + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); - 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); - } - } + // CANCELLED status should be excluded from call counting, so no ejections occur. + assertEjectedSubchannels(ImmutableSet.of()); + } + /** + * The success rate algorithm ignores DEADLINE_EXCEEDED status calls. + */ + @Test + public void successRateOneOutlier_deadlineExceededIgnored() { + 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)); + + // subchannel1 returns DEADLINE_EXCEEDED. + generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + + // Move forward in time to a point where the detection timer has fired. forwardTime(config); - // subchannel1 was cancelled client-side and should not be ejected as an outlier. + // DEADLINE_EXCEEDED status should be excluded from call counting, so no ejections occur. assertEjectedSubchannels(ImmutableSet.of()); } @@ -615,7 +620,7 @@ public void successRateOneOutlier_configChange() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -635,7 +640,7 @@ public void successRateOneOutlier_configChange() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel2, Status.DEADLINE_EXCEEDED), 8); + generateLoad(ImmutableMap.of(subchannel2, Status.UNAVAILABLE), 8); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -660,7 +665,7 @@ public void successRateOneOutlier_unejected() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. fakeClock.forwardTime(config.intervalNanos + 1, TimeUnit.NANOSECONDS); @@ -695,7 +700,7 @@ public void successRateOneOutlier_notEnoughVolume() { // We produce an outlier, but don't give it enough calls to reach the minimum volume. generateLoad( - ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), + ImmutableMap.of(subchannel1, Status.UNAVAILABLE), ImmutableMap.of(subchannel1, 19), 7); // Move forward in time to a point where the detection timer has fired. @@ -722,7 +727,7 @@ public void successRateOneOutlier_notEnoughAddressesWithVolume() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); generateLoad( - ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), + ImmutableMap.of(subchannel1, Status.UNAVAILABLE), // subchannel2 has only 19 calls which results in success rate not triggering. ImmutableMap.of(subchannel2, 19), 7); @@ -751,7 +756,7 @@ public void successRateOneOutlier_enforcementPercentage() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -777,8 +782,8 @@ public void successRateTwoOutliers() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); generateLoad(ImmutableMap.of( - subchannel1, Status.DEADLINE_EXCEEDED, - subchannel2, Status.DEADLINE_EXCEEDED), 7); + subchannel1, Status.UNAVAILABLE, + subchannel2, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -806,9 +811,9 @@ public void successRateThreeOutliers_maxEjectionPercentage() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); generateLoad(ImmutableMap.of( - subchannel1, Status.DEADLINE_EXCEEDED, - subchannel2, Status.DEADLINE_EXCEEDED, - subchannel3, Status.DEADLINE_EXCEEDED), 7); + subchannel1, Status.UNAVAILABLE, + subchannel2, Status.UNAVAILABLE, + subchannel3, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -910,7 +915,7 @@ public void failurePercentageOneOutlier() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -919,6 +924,56 @@ public void failurePercentageOneOutlier() { assertEjectedSubchannels(ImmutableSet.of(ImmutableSet.copyOf(servers.get(0).getAddresses()))); } + /** + * The failure percentage algorithm ignores CANCELLED status calls. + */ + @Test + public void failurePercentageOneOutlier_cancelledIgnored() { + 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)); + + // subchannel1 returns CANCELLED. + generateLoad(ImmutableMap.of(subchannel1, Status.CANCELLED), 7); + + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); + + // CANCELLED status should be excluded from call counting, so no ejections occur. + assertEjectedSubchannels(ImmutableSet.of()); + } + + /** + * The failure percentage algorithm ignores DEADLINE_EXCEEDED status calls. + */ + @Test + public void failurePercentageOneOutlier_deadlineExceededIgnored() { + 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)); + + // subchannel1 returns DEADLINE_EXCEEDED. + generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); + + // DEADLINE_EXCEEDED status should be excluded from call counting, so no ejections occur. + assertEjectedSubchannels(ImmutableSet.of()); + } + /** * The failure percentage algorithm ignores addresses without enough volume.. */ @@ -934,7 +989,7 @@ public void failurePercentageOneOutlier_notEnoughVolume() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -960,7 +1015,7 @@ public void failurePercentageOneOutlier_notEnoughAddressesWithVolume() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); generateLoad( - ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), + ImmutableMap.of(subchannel1, Status.UNAVAILABLE), // subchannel2 has only 19 calls which results in failure percentage not triggering. ImmutableMap.of(subchannel2, 19), 7); @@ -989,7 +1044,7 @@ public void failurePercentageOneOutlier_enforcementPercentage() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1023,9 +1078,9 @@ public void successRateAndFailurePercentageThreeOutliers() { // configured with a 0 tolerance threshold. generateLoad( ImmutableMap.of( - subchannel1, Status.DEADLINE_EXCEEDED, - subchannel2, Status.DEADLINE_EXCEEDED, - subchannel3, Status.DEADLINE_EXCEEDED), + subchannel1, Status.UNAVAILABLE, + subchannel2, Status.UNAVAILABLE, + subchannel3, Status.UNAVAILABLE), ImmutableMap.of(subchannel3, 1), 7); // Move forward in time to a point where the detection timer has fired. @@ -1055,7 +1110,7 @@ public void subchannelUpdateAddress_singleReplaced() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1117,8 +1172,8 @@ public void multipleAddressesEndpoint() { assertThat(loadBalancer.endpointTrackerMap.size()).isEqualTo(3); assertThat(loadBalancer.addressMap.size()).isEqualTo(5); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED, - subchannel2, Status.DEADLINE_EXCEEDED), 13); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE, + subchannel2, Status.UNAVAILABLE), 13); forwardTime(config); // eject the first endpoint: (address0, address1) @@ -1186,7 +1241,7 @@ public void subchannelUpdateAddress_multipleReplacedWithSingle() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 6); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 6); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1276,7 +1331,7 @@ public void successRateAndFailurePercentage_successRateOutlier() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1305,7 +1360,7 @@ public void successRateAndFailurePercentage_successRateOutlier_() { // with heal loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 6); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 6); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1349,7 +1404,7 @@ public void successRateAndFailurePercentage_errorPercentageOutlier() { loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 7); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1378,7 +1433,7 @@ public void successRateAndFailurePercentage_errorPercentageOutlier_() { // with loadBalancer.acceptResolvedAddresses(buildResolvedAddress(config, servers)); - generateLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 6); + generateLoad(ImmutableMap.of(subchannel1, Status.UNAVAILABLE), 6); // Move forward in time to a point where the detection timer has fired. forwardTime(config); @@ -1477,8 +1532,12 @@ private void generateLoad(Map statusMap, int calls = callCountMap.containsKey(subchannel) ? callCountMap.get(subchannel) : 0; if (calls < maxCalls) { callCountMap.put(subchannel, ++calls); - clientStreamTracer.streamClosed( - statusMap.containsKey(subchannel) ? statusMap.get(subchannel) : Status.OK); + Status status = statusMap.containsKey(subchannel) ? statusMap.get(subchannel) : Status.OK; + if (status.getCode() == Status.Code.CANCELLED + || status.getCode() == Status.Code.DEADLINE_EXCEEDED) { + clientStreamTracer.cancelled(status); + } + clientStreamTracer.streamClosed(status); } } } From 6d3cf19e2107ee235c0f62b2d26262684ee13dd3 Mon Sep 17 00:00:00 2001 From: agrawalabhi Date: Mon, 3 Aug 2026 12:35:41 +0000 Subject: [PATCH 11/15] core: trigger clientCancelled in closeListener based on stopDelivery to avoid races --- .../java/io/grpc/internal/AbstractClientStream.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/core/src/main/java/io/grpc/internal/AbstractClientStream.java b/core/src/main/java/io/grpc/internal/AbstractClientStream.java index b3995b5380c..2e5b333bc36 100644 --- a/core/src/main/java/io/grpc/internal/AbstractClientStream.java +++ b/core/src/main/java/io/grpc/internal/AbstractClientStream.java @@ -201,7 +201,6 @@ public final void cancel(Status reason) { return; } cancelled = true; - transportState().getStatsTraceContext().clientCancelled(reason); abstractClientStreamSink().cancel(reason); } @@ -443,13 +442,13 @@ public final void transportReportStatus( if (deframerClosed) { deframerClosedTask = null; - closeListener(status, rpcProgress, trailers); + closeListener(status, rpcProgress, trailers, stopDelivery); } else { deframerClosedTask = new Runnable() { @Override public void run() { - closeListener(status, rpcProgress, trailers); + closeListener(status, rpcProgress, trailers, stopDelivery); } }; closeDeframer(stopDelivery); @@ -462,9 +461,12 @@ public void run() { * @throws IllegalStateException if the call has not yet been started. */ private void closeListener( - Status status, RpcProgress rpcProgress, Metadata trailers) { + Status status, RpcProgress rpcProgress, Metadata trailers, boolean stopDelivery) { if (!listenerClosed) { listenerClosed = true; + if (stopDelivery) { + statsTraceCtx.clientCancelled(status); + } statsTraceCtx.streamClosed(status); if (getTransportTracer() != null) { getTransportTracer().reportStreamClosed(status.isOk()); From 9331806f4265c6b9d6e1bfd842f6393c9a82699f Mon Sep 17 00:00:00 2001 From: agrawalabhi Date: Mon, 3 Aug 2026 12:38:29 +0000 Subject: [PATCH 12/15] util: restore incrementCallCount(boolean success) method signature --- .../java/io/grpc/util/OutlierDetectionLoadBalancer.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java b/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java index 966affc6cff..b7fd50ae6cd 100644 --- a/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java +++ b/util/src/main/java/io/grpc/util/OutlierDetectionLoadBalancer.java @@ -493,7 +493,7 @@ public void cancelled(Status status) { @Override public void streamClosed(Status status) { if (!cancelled) { - tracker.incrementCallCount(status); + tracker.incrementCallCount(status.isOk()); } delegate().streamClosed(status); } @@ -510,7 +510,7 @@ public void cancelled(Status status) { @Override public void streamClosed(Status status) { if (!cancelled) { - tracker.incrementCallCount(status); + tracker.incrementCallCount(status.isOk()); } } }; @@ -571,13 +571,13 @@ Set getSubchannels() { return ImmutableSet.copyOf(subchannels); } - void incrementCallCount(Status status) { + void incrementCallCount(boolean success) { // If neither algorithm is configured, no point in incrementing counters. if (config.successRateEjection == null && config.failurePercentageEjection == null) { return; } - if (status.isOk()) { + if (success) { activeCallCounter.successCount.getAndIncrement(); } else { activeCallCounter.failureCount.getAndIncrement(); From 5f3af3d86ee5fbba2800d8d1604c3b561a788744 Mon Sep 17 00:00:00 2001 From: agrawalabhi Date: Mon, 3 Aug 2026 12:52:14 +0000 Subject: [PATCH 13/15] test: add explicit stopDelivery unit tests in AbstractClientStreamTest --- .../internal/AbstractClientStreamTest.java | 33 ++++++++++++++++++- 1 file changed, 32 insertions(+), 1 deletion(-) diff --git a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java index 346b6fe6e15..29d722a193f 100644 --- a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java @@ -192,9 +192,40 @@ public void cancel(Status errorStatus) { state.transportReportStatus(Status.OK, false, new Metadata()); // Subsequent late cancellation - stream.cancel(Status.CANCELLED.withDescription("Late cancel")); + verify(mockTracer, never()).cancelled(any(Status.class)); + } + + @Test + public void transportReportStatus_stopDeliveryFalse_doesNotNotifyTracerCancelled() { + 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() {}, + customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + // Server-initiated CANCELLED (stopDelivery = false) + state.transportReportStatus(Status.CANCELLED, false, new Metadata()); verify(mockTracer, never()).cancelled(any(Status.class)); + verify(mockTracer).streamClosed(Status.CANCELLED); + } + + @Test + public void transportReportStatus_stopDeliveryTrue_notifiesTracerCancelled() { + 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() {}, + customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + // Client/Transport-initiated cancellation (stopDelivery = true) + Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled"); + state.transportReportStatus(cancelStatus, true, new Metadata()); + + verify(mockTracer).cancelled(cancelStatus); + verify(mockTracer).streamClosed(cancelStatus); } @Test From be50f7b7da36894ed536f8a590f46d9902ec54ee Mon Sep 17 00:00:00 2001 From: agrawalabhi Date: Mon, 3 Aug 2026 12:56:21 +0000 Subject: [PATCH 14/15] test: add comprehensive stopDelivery and server vs client cancellation unit tests --- .../internal/AbstractClientStreamTest.java | 90 ++++++++++++ .../inprocess/InProcessTransportTest.java | 88 ++++++++++++ .../OutlierDetectionLoadBalancerTest.java | 134 ++++++++++++++++++ 3 files changed, 312 insertions(+) diff --git a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java index 29d722a193f..6c24e9e5b0f 100644 --- a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java @@ -228,6 +228,96 @@ public void transportReportStatus_stopDeliveryTrue_notifiesTracerCancelled() { verify(mockTracer).streamClosed(cancelStatus); } + @Test + public void transportReportStatus_stopDeliveryFalse_deadlineExceeded_noTracerCancelled() { + 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() {}, customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + // Server-initiated DEADLINE_EXCEEDED (stopDelivery = false) + Status status = Status.DEADLINE_EXCEEDED.withDescription("Server deadline exceeded"); + state.transportReportStatus(status, false, new Metadata()); + + verify(mockTracer, never()).cancelled(any(Status.class)); + verify(mockTracer).streamClosed(status); + } + + @Test + public void transportReportStatus_stopDeliveryTrue_deadlineExceeded_notifiesTracerCancelled() { + 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() {}, customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + // Client/Transport-initiated deadline exceeded (stopDelivery = true) + Status status = Status.DEADLINE_EXCEEDED.withDescription("Client deadline exceeded"); + state.transportReportStatus(status, true, new Metadata()); + + verify(mockTracer).cancelled(status); + verify(mockTracer).streamClosed(status); + } + + @Test + public void closeListener_directAssertions_stopDeliveryTrueAndFalse() { + ClientStreamTracer mockTracer1 = mock(ClientStreamTracer.class); + StatsTraceContext statsTraceCtx1 = new StatsTraceContext(new StreamTracer[] {mockTracer1}); + BaseTransportState state1 = new BaseTransportState(statsTraceCtx1, transportTracer); + AbstractClientStream stream1 = new BaseAbstractClientStream( + allocator, state1, new BaseSink() {}, statsTraceCtx1, transportTracer); + stream1.start(mockListener); + + // stopDelivery = true: clientCancelled is called before streamClosed + Status statusTrue = Status.CANCELLED.withDescription("stopDelivery true"); + state1.transportReportStatus(statusTrue, true, new Metadata()); + verify(mockTracer1).cancelled(statusTrue); + verify(mockTracer1).streamClosed(statusTrue); + + ClientStreamTracer mockTracer2 = mock(ClientStreamTracer.class); + StatsTraceContext statsTraceCtx2 = new StatsTraceContext(new StreamTracer[] {mockTracer2}); + BaseTransportState state2 = new BaseTransportState(statsTraceCtx2, transportTracer); + AbstractClientStream stream2 = new BaseAbstractClientStream( + allocator, state2, new BaseSink() {}, statsTraceCtx2, transportTracer); + stream2.start(mockListener); + + // stopDelivery = false: clientCancelled is NOT called, only streamClosed + Status statusFalse = Status.CANCELLED.withDescription("stopDelivery false"); + state2.transportReportStatus(statusFalse, false, new Metadata()); + verify(mockTracer2, never()).cancelled(any(Status.class)); + verify(mockTracer2).streamClosed(statusFalse); + } + + @Test + public void closeListener_deferredDeframerClose_stopDeliveryFalse_delaysCloseListener() { + ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); + StatsTraceContext customStatsTraceCtx = new StatsTraceContext(new StreamTracer[] {mockTracer}); + BaseTransportState state = new BaseTransportState(customStatsTraceCtx, transportTracer); + AbstractClientStream stream = new BaseAbstractClientStream( + allocator, state, new BaseSink() {}, customStatsTraceCtx, transportTracer); + stream.start(mockListener); + + // Send partial message into deframer + byte[] data = new byte[] {0, 0, 0, 0, 2, 1}; // 2-byte frame, only 1 byte delivered + state.deframe(ReadableBuffers.wrap(data)); + + Status statusFalse = Status.CANCELLED.withDescription("deferred stopDelivery false"); + state.transportReportStatus(statusFalse, false, new Metadata()); + + // Listener is not closed yet because deframer is mid-frame and waiting for complete frame + verify(mockTracer, never()).cancelled(any(Status.class)); + verify(mockTracer, never()).streamClosed(any(Status.class)); + + // Request message and provide remaining byte of frame to complete deframer processing + stream.request(1); + state.deframe(ReadableBuffers.wrap(new byte[] {2})); + verify(mockTracer, never()).cancelled(any(Status.class)); + verify(mockTracer).streamClosed(any(Status.class)); + } + @Test public void startFailsOnNullListener() { AbstractClientStream stream = diff --git a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java index 1f977eaecb5..6955e749469 100644 --- a/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java +++ b/inprocess/src/test/java/io/grpc/inprocess/InProcessTransportTest.java @@ -287,6 +287,94 @@ methodDescriptor, new Metadata(), CallOptions.DEFAULT, .cancelled(org.mockito.Mockito.any(Status.class)); } + @Test + public void serverStream_closeWithCancelled_doesNotNotifyTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); + server.start(serverListener); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener = + serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + ClientStreamTracer mockTracer = org.mockito.Mockito.mock(ClientStreamTracer.class); + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), CallOptions.DEFAULT, + new ClientStreamTracer[] {mockTracer}); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + StreamCreation serverStreamCreation = + serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + ServerStream serverStream = serverStreamCreation.stream; + + Status serverStatus = Status.CANCELLED.withDescription("Server cancelled over wire"); + serverStream.close(serverStatus, new Metadata()); + + org.mockito.Mockito.verify(mockTracer, org.mockito.Mockito.never()) + .cancelled(org.mockito.Mockito.any(Status.class)); + org.mockito.ArgumentCaptor statusCaptor = + org.mockito.ArgumentCaptor.forClass(Status.class); + org.mockito.Mockito.verify(mockTracer).streamClosed(statusCaptor.capture()); + assertEquals(Status.Code.CANCELLED, statusCaptor.getValue().getCode()); + assertEquals("Server cancelled over wire", statusCaptor.getValue().getDescription()); + } + + @Test + public void serverStream_closeWithDeadlineExceeded_noTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); + server.start(serverListener); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener = + serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + ClientStreamTracer mockTracer = org.mockito.Mockito.mock(ClientStreamTracer.class); + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), CallOptions.DEFAULT, + new ClientStreamTracer[] {mockTracer}); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + StreamCreation serverStreamCreation = + serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + ServerStream serverStream = serverStreamCreation.stream; + + Status serverStatus = Status.DEADLINE_EXCEEDED.withDescription("Server deadline exceeded"); + serverStream.close(serverStatus, new Metadata()); + + org.mockito.Mockito.verify(mockTracer, org.mockito.Mockito.never()) + .cancelled(org.mockito.Mockito.any(Status.class)); + org.mockito.ArgumentCaptor statusCaptor = + org.mockito.ArgumentCaptor.forClass(Status.class); + org.mockito.Mockito.verify(mockTracer).streamClosed(statusCaptor.capture()); + assertEquals(Status.Code.DEADLINE_EXCEEDED, statusCaptor.getValue().getCode()); + assertEquals("Server deadline exceeded", statusCaptor.getValue().getDescription()); + } + + @Test + public void clientStream_cancelWithDeadlineExceeded_notifiesTracerCancelled() throws Exception { + server = newServer(Arrays.asList(serverStreamTracerFactory)); + server.start(serverListener); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener = + serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + ClientStreamTracer mockTracer = org.mockito.Mockito.mock(ClientStreamTracer.class); + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), CallOptions.DEFAULT, + new ClientStreamTracer[] {mockTracer}); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + + Status cancelStatus = Status.DEADLINE_EXCEEDED.withDescription("Client deadline exceeded"); + clientStream.cancel(cancelStatus); + + org.mockito.Mockito.verify(mockTracer).cancelled(cancelStatus); + org.mockito.Mockito.verify(mockTracer).streamClosed(cancelStatus); + } + private void assertAssumedMessageSize( TestStreamTracer streamTracerSender, TestStreamTracer streamTracerReceiver) { if (isEnabledSupportTracingMessageSizes()) { diff --git a/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java b/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java index 0e008d3dfce..2e1aabfb5a9 100644 --- a/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java +++ b/util/src/test/java/io/grpc/util/OutlierDetectionLoadBalancerTest.java @@ -604,6 +604,58 @@ public void successRateOneOutlier_deadlineExceededIgnored() { assertEjectedSubchannels(ImmutableSet.of()); } + /** + * Server-initiated CANCELLED status over the wire (stopDelivery = false) counts as failure + * and results in ejection under success rate algorithm. + */ + @Test + public void successRateOneOutlier_serverInitiatedCancelledEjected() { + 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)); + + // subchannel1 returns CANCELLED from server (no client cancellation tracer callback). + generateServerInitiatedLoad(ImmutableMap.of(subchannel1, Status.CANCELLED), 7); + + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); + + // Server-initiated CANCELLED status is counted as a failure, so subchannel1 should be ejected. + assertEjectedSubchannels(ImmutableSet.of(ImmutableSet.copyOf(servers.get(0).getAddresses()))); + } + + /** + * Server-initiated DEADLINE_EXCEEDED status over the wire (stopDelivery = false) counts + * as failure and results in ejection under success rate algorithm. + */ + @Test + public void successRateOneOutlier_serverInitiatedDeadlineExceededEjected() { + 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)); + + // subchannel1 returns DEADLINE_EXCEEDED from server (no client cancellation tracer callback). + generateServerInitiatedLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); + + // Server-initiated DEADLINE_EXCEEDED status is counted as a failure, so subchannel1 is ejected. + assertEjectedSubchannels(ImmutableSet.of(ImmutableSet.copyOf(servers.get(0).getAddresses()))); + } + /** * The success rate algorithm ejects the outlier, but then the config changes so that similar * behavior no longer gets ejected. @@ -974,6 +1026,58 @@ public void failurePercentageOneOutlier_deadlineExceededIgnored() { assertEjectedSubchannels(ImmutableSet.of()); } + /** + * Server-initiated CANCELLED status over the wire (stopDelivery = false) counts as failure + * and results in ejection under failure percentage algorithm. + */ + @Test + public void failurePercentageOneOutlier_serverInitiatedCancelledEjected() { + 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)); + + // subchannel1 returns CANCELLED from server (no client cancellation tracer callback). + generateServerInitiatedLoad(ImmutableMap.of(subchannel1, Status.CANCELLED), 7); + + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); + + // Server-initiated CANCELLED status is counted as a failure, so subchannel1 should be ejected. + assertEjectedSubchannels(ImmutableSet.of(ImmutableSet.copyOf(servers.get(0).getAddresses()))); + } + + /** + * Server-initiated DEADLINE_EXCEEDED status over the wire (stopDelivery = false) counts + * as failure and results in ejection under failure percentage algorithm. + */ + @Test + public void failurePercentageOneOutlier_serverInitiatedDeadlineExceededEjected() { + 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)); + + // subchannel1 returns DEADLINE_EXCEEDED from server (no client cancellation tracer callback). + generateServerInitiatedLoad(ImmutableMap.of(subchannel1, Status.DEADLINE_EXCEEDED), 7); + + // Move forward in time to a point where the detection timer has fired. + forwardTime(config); + + // Server-initiated DEADLINE_EXCEEDED status is counted as a failure, so subchannel1 is ejected. + assertEjectedSubchannels(ImmutableSet.of(ImmutableSet.copyOf(servers.get(0).getAddresses()))); + } + /** * The failure percentage algorithm ignores addresses without enough volume.. */ @@ -1542,6 +1646,36 @@ private void generateLoad(Map statusMap, } } + // Generates 100 calls, simulating server-initiated status responses over the wire. + private void generateServerInitiatedLoad( + Map statusMap, int expectedStateChanges) { + 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(expectedStateChanges)).updateBalancingState(stateCaptor.capture(), + pickerCaptor.capture()); + SubchannelPicker picker = pickerCaptor.getAllValues() + .get(pickerCaptor.getAllValues().size() - 1); + + HashMap callCountMap = new HashMap<>(); + 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(); + + int calls = callCountMap.containsKey(subchannel) ? callCountMap.get(subchannel) : 0; + callCountMap.put(subchannel, ++calls); + Status status = statusMap.containsKey(subchannel) ? statusMap.get(subchannel) : Status.OK; + clientStreamTracer.streamClosed(status); + } + } + // Forwards time past the moment when the timer will fire. private void forwardTime(OutlierDetectionLoadBalancerConfig config) { fakeClock.forwardTime(config.intervalNanos + 1, TimeUnit.NANOSECONDS); From c6e2f4aeef762761fd15cf8e2e25b2291e27fa04 Mon Sep 17 00:00:00 2001 From: agrawalabhi Date: Mon, 3 Aug 2026 12:59:35 +0000 Subject: [PATCH 15/15] test: remove redundant directAssertions test case --- .../internal/AbstractClientStreamTest.java | 29 ------------------- 1 file changed, 29 deletions(-) diff --git a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java index 6c24e9e5b0f..ddcf1bf44f4 100644 --- a/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java @@ -262,35 +262,6 @@ public void transportReportStatus_stopDeliveryTrue_deadlineExceeded_notifiesTrac verify(mockTracer).streamClosed(status); } - @Test - public void closeListener_directAssertions_stopDeliveryTrueAndFalse() { - ClientStreamTracer mockTracer1 = mock(ClientStreamTracer.class); - StatsTraceContext statsTraceCtx1 = new StatsTraceContext(new StreamTracer[] {mockTracer1}); - BaseTransportState state1 = new BaseTransportState(statsTraceCtx1, transportTracer); - AbstractClientStream stream1 = new BaseAbstractClientStream( - allocator, state1, new BaseSink() {}, statsTraceCtx1, transportTracer); - stream1.start(mockListener); - - // stopDelivery = true: clientCancelled is called before streamClosed - Status statusTrue = Status.CANCELLED.withDescription("stopDelivery true"); - state1.transportReportStatus(statusTrue, true, new Metadata()); - verify(mockTracer1).cancelled(statusTrue); - verify(mockTracer1).streamClosed(statusTrue); - - ClientStreamTracer mockTracer2 = mock(ClientStreamTracer.class); - StatsTraceContext statsTraceCtx2 = new StatsTraceContext(new StreamTracer[] {mockTracer2}); - BaseTransportState state2 = new BaseTransportState(statsTraceCtx2, transportTracer); - AbstractClientStream stream2 = new BaseAbstractClientStream( - allocator, state2, new BaseSink() {}, statsTraceCtx2, transportTracer); - stream2.start(mockListener); - - // stopDelivery = false: clientCancelled is NOT called, only streamClosed - Status statusFalse = Status.CANCELLED.withDescription("stopDelivery false"); - state2.transportReportStatus(statusFalse, false, new Metadata()); - verify(mockTracer2, never()).cancelled(any(Status.class)); - verify(mockTracer2).streamClosed(statusFalse); - } - @Test public void closeListener_deferredDeframerClose_stopDeliveryFalse_delaysCloseListener() { ClientStreamTracer mockTracer = mock(ClientStreamTracer.class);