Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
9d14356
outlier-detection: exclude client and hedging cancellations from call…
AgraVator Jul 22, 2026
7abb485
Update @since tag to 1.84.0 and add tracer cancellation in FailingCli…
AgraVator Jul 22, 2026
f9cc87f
core, inprocess, binder: synchronize client cancellation tracer notif…
AgraVator Jul 23, 2026
120e35c
test: add unit tests for race condition and transport tracer cancella…
AgraVator Jul 24, 2026
f9dd0fb
test: add inprocess client stream cancellation unit tests
AgraVator Jul 24, 2026
fe581ae
Revert "test: add inprocess client stream cancellation unit tests"
AgraVator Jul 24, 2026
5e558a5
test: fix checkstyle line length violations in InProcessTransportTest
AgraVator Jul 24, 2026
634be81
Revert "test: fix checkstyle line length violations in InProcessTrans…
AgraVator Jul 24, 2026
20682d9
test: fix checkstyle violations and missing server.start in InProcess…
AgraVator Jul 27, 2026
bf90eb5
util: exclude client cancellations and deadline exceeded from outlier…
AgraVator Aug 3, 2026
6d3cf19
core: trigger clientCancelled in closeListener based on stopDelivery …
AgraVator Aug 3, 2026
9331806
util: restore incrementCallCount(boolean success) method signature
AgraVator Aug 3, 2026
5f3af3d
test: add explicit stopDelivery unit tests in AbstractClientStreamTest
AgraVator Aug 3, 2026
be50f7b
test: add comprehensive stopDelivery and server vs client cancellatio…
AgraVator Aug 3, 2026
c6e2f4a
test: remove redundant directAssertions test case
AgraVator Aug 3, 2026
f27ce70
core, binder: notify clientCancelled directly in transportReportStatu…
AgraVator Aug 5, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions api/src/main/java/io/grpc/ClientStreamTracer.java
Original file line number Diff line number Diff line change
Expand Up @@ -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.84.0
*/
public void cancelled(Status status) {
}

/**
* Factory class for {@link ClientStreamTracer}.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -292,6 +292,22 @@ public void testMessageProducerClosedAfterStream_b169313545() throws Exception {
streamListener.drainMessages();
}

@Test
public void testCancelStream_notifiesTracer() throws Exception {
transport = new BinderClientTransportBuilder().build();
startAndAwaitReady(transport, transportListener);

ClientStreamTracer mockTracer = org.mockito.Mockito.mock(ClientStreamTracer.class);
ClientStream stream =
transport.newStream(methodDesc, new Metadata(), CallOptions.DEFAULT, new ClientStreamTracer[]{mockTracer});

stream.start(streamListener);
Status cancelStatus = Status.CANCELLED.withDescription("Client cancelled");
stream.cancel(cancelStatus);

org.mockito.Mockito.verify(mockTracer, org.mockito.Mockito.timeout(5000)).cancelled(org.mockito.ArgumentMatchers.eq(cancelStatus));
}

@Test
public void testNewStreamBeforeTransportReadyFails() throws Exception {
// Use a special SecurityPolicy that lets us act before the transport is setup/ready.
Expand Down
3 changes: 3 additions & 0 deletions binder/src/main/java/io/grpc/binder/internal/Inbound.java
Original file line number Diff line number Diff line change
Expand Up @@ -268,6 +268,9 @@ private final void deliverInternal() {

@GuardedBy("this")
final void closeOnCancel(Status status) {
if (!isClosed() && statsTraceContext != null) {
Comment thread
AgraVator marked this conversation as resolved.
statsTraceContext.clientCancelled(status);
}
closeAbnormal(Status.CANCELLED, status, false);
}

Expand Down
21 changes: 13 additions & 8 deletions core/src/main/java/io/grpc/internal/AbstractClientStream.java
Original file line number Diff line number Diff line change
Expand Up @@ -251,6 +251,8 @@ protected TransportState(
}
}



private void setFullStreamDecompression(boolean fullStreamDecompression) {
this.fullStreamDecompression = fullStreamDecompression;
}
Expand Down Expand Up @@ -391,16 +393,16 @@ protected void inboundTrailersReceived(Metadata trailers, Status status) {
* method must be called from the transport thread.
*
* @param status the new status to set
* @param stopDelivery if {@code true}, interrupts any further delivery of inbound messages that
* @param cancelled if {@code true}, interrupts any further delivery of inbound messages that
* may already be queued up in the deframer. If {@code false}, the listener will be
* notified immediately after all currently completed messages in the deframer have been
* delivered to the application.
* @param trailers new instance of {@code Trailers}, either empty or those returned by the
* server
*/
public final void transportReportStatus(final Status status, boolean stopDelivery,
public final void transportReportStatus(final Status status, boolean cancelled,
final Metadata trailers) {
transportReportStatus(status, RpcProgress.PROCESSED, stopDelivery, trailers);
transportReportStatus(status, RpcProgress.PROCESSED, cancelled, trailers);
}

/**
Expand All @@ -411,7 +413,7 @@ public final void transportReportStatus(final Status status, boolean stopDeliver
* @param rpcProgress RPC progress that the
* {@link ClientStreamListener#closed(Status, RpcProgress, Metadata)}
* will receive
* @param stopDelivery if {@code true}, interrupts any further delivery of inbound messages that
* @param cancelled if {@code true}, interrupts any further delivery of inbound messages that
* may already be queued up in the deframer and overrides any previously queued status.
* If {@code false}, the listener will be notified immediately after all currently
* completed messages in the deframer have been delivered to the application.
Expand All @@ -421,16 +423,19 @@ public final void transportReportStatus(final Status status, boolean stopDeliver
public final void transportReportStatus(
final Status status,
final RpcProgress rpcProgress,
boolean stopDelivery,
boolean cancelled,
final Metadata trailers) {
checkNotNull(status, "status");
checkNotNull(trailers, "trailers");
// If stopDelivery, we continue in case previous invocation is waiting for stall
if (statusReported && !stopDelivery) {
// If cancelled, we continue in case previous invocation is waiting for stall
if (statusReported && !cancelled) {
return;
}
statusReported = true;
statusReportedIsOk = status.isOk();
if (cancelled) {
statsTraceCtx.clientCancelled(status);
}
onStreamDeallocated();

if (deframerClosed) {
Expand All @@ -444,7 +449,7 @@ public void run() {
closeListener(status, rpcProgress, trailers);
}
};
closeDeframer(stopDelivery);
closeDeframer(cancelled);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
13 changes: 13 additions & 0 deletions core/src/main/java/io/grpc/internal/StatsTraceContext.java
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,19 @@ public void serverCallMethodResolved(MethodDescriptor<?, ?> method) {
}
}

/**
* See {@link ClientStreamTracer#cancelled}. For client-side only.
*
* <p>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.
Expand Down
134 changes: 134 additions & 0 deletions core/src/test/java/io/grpc/internal/AbstractClientStreamTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -155,6 +156,139 @@ 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) {
state.transportReportStatus(errorStatus, true, new Metadata());
}
}, customStatsTraceCtx, transportTracer);
stream.start(mockListener);

Status cancelStatus = Status.CANCELLED.withDescription("Cancelled by test");
stream.cancel(cancelStatus);

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
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
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_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 =
Expand Down
49 changes: 49 additions & 0 deletions core/src/test/java/io/grpc/internal/StatsTraceContextTest.java
Original file line number Diff line number Diff line change
@@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -846,6 +846,7 @@ public void cancel(Status reason) {
if (!internalCancel(serverStatus, serverStatus)) {
return;
}
statsTraceCtx.clientCancelled(reason);
serverStream.clientCancelled(reason);
streamClosed();
}
Expand Down
Loading
Loading