Skip to content

Commit 6a55ff2

Browse files
committed
add missing endDelay()
1 parent a992bdf commit 6a55ff2

2 files changed

Lines changed: 27 additions & 2 deletions

File tree

core/src/main/java/io/grpc/internal/DelayedClientTransport.java

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -247,7 +247,7 @@ public final void shutdownNow(Status status) {
247247
}
248248
if (savedReportTransportTerminated != null) {
249249
for (PendingStream stream : savedPendingStreams) {
250-
Runnable runnable = stream.setStream(
250+
Runnable runnable = stream.setStreamAndEndDelay(
251251
new FailingClientStream(status, RpcProgress.REFUSED, stream.tracers));
252252
if (runnable != null) {
253253
// Drain in-line instead of using an executor as failing stream just throws everything
@@ -406,6 +406,11 @@ void endDelay() {
406406
}
407407
}
408408

409+
Runnable setStreamAndEndDelay(ClientStream stream) {
410+
endDelay();
411+
return setStream(stream);
412+
}
413+
409414
/** Runnable may be null. */
410415
private Runnable createRealStream(ClientTransport transport, String authorityOverride) {
411416
ClientStream realStream;
@@ -424,7 +429,7 @@ private Runnable createRealStream(ClientTransport transport, String authorityOve
424429
// been called on the delayed stream.
425430
realStream.setAuthority(authorityOverride);
426431
}
427-
return setStream(realStream);
432+
return setStreamAndEndDelay(realStream);
428433
}
429434

430435
@Override

core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -821,6 +821,26 @@ public void streamDelayMetrics_cancelled() {
821821
verify(mockTracer).delayEnded();
822822
}
823823

824+
@Test
825+
public void streamDelayMetrics_shutdownNow() {
826+
ClientStreamTracer mockTracer = mock(ClientStreamTracer.class);
827+
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer };
828+
829+
SubchannelPicker connectingPicker = mock(SubchannelPicker.class);
830+
when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class)))
831+
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
832+
833+
delayedTransport.reprocess(connectingPicker);
834+
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
835+
stream.start(streamListener);
836+
837+
verify(mockTracer).delayStarted("pick_first:connecting");
838+
839+
delayedTransport.shutdownNow(Status.UNAVAILABLE);
840+
841+
verify(mockTracer).delayEnded();
842+
}
843+
824844
private static TransportProvider newTransportProvider(final ClientTransport transport) {
825845
return new TransportProvider() {
826846
@Override

0 commit comments

Comments
 (0)