From 9d14356cea87c2cc8cd44172ea3c2dab3074832f Mon Sep 17 00:00:00 2001 From: AgraVator Date: Wed, 22 Jul 2026 20:01:28 +0530 Subject: [PATCH 1/8] 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/8] 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 3/8] 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 4/8] 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 5/8] 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 6/8] 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 7/8] 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 8/8] 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(