Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
13 changes: 12 additions & 1 deletion core/src/main/java/io/grpc/internal/RetriableStream.java
Original file line number Diff line number Diff line change
Expand Up @@ -399,8 +399,19 @@ public final void start(ClientStreamListener listener) {
return;
}

// cancel() may have committed this stream before prestart() registered it. In that case the
// post-commit callback ran before registration and could not remove the stream. Run it again
// after registration so the channel's uncommitted stream registry cannot retain this stream.
boolean alreadyCommitted;
synchronized (lock) {
state.buffer.add(new StartEntry());
alreadyCommitted = state.winningSubstream != null;
if (!alreadyCommitted) {
state.buffer.add(new StartEntry());
}
}
if (alreadyCommitted) {
postCommit();
return;
}

Substream substream = createSubstream(0, false, false);
Expand Down
40 changes: 40 additions & 0 deletions core/src/test/java/io/grpc/internal/RetriableStreamTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,46 @@ public void tearDown() {
assertEquals(0, fakeClock.numPendingTasks());
}

@Test
public void cancelDuringStartBeforeRegistration() {
assertCancelDuringStart(retriableStream);
}

@Test
public void hedging_cancelBeforeRegistration() {
assertCancelDuringStart(hedgingStream);
}

@Test
public void transparentRetry_cancelBeforeRegistration() {
RetriableStream<String> stream = new RecordedRetriableStream(
method, new Metadata(), channelBufferUsed, PER_RPC_BUFFER_LIMIT, CHANNEL_BUFFER_LIMIT,
MoreExecutors.directExecutor(), fakeClock.getScheduledExecutorService(), null, null, null);
assertCancelDuringStart(stream);
}

private void assertCancelDuringStart(RetriableStream<String> stream) {
Status reason = Status.CANCELLED.withDescription("cancel before registration");
List<RetriableStream<String>> registeredStreams = new ArrayList<>();
doAnswer(invocation -> {
stream.cancel(reason);
registeredStreams.add(stream);
return null;
}).when(retriableStreamRecorder).prestart();
doAnswer(invocation -> {
registeredStreams.remove(stream);
return null;
}).when(retriableStreamRecorder).postCommit();

stream.start(masterListener);

assertThat(registeredStreams).isEmpty();
verify(retriableStreamRecorder).prestart();
verify(retriableStreamRecorder, times(2)).postCommit();
verify(masterListener).closed(same(reason), same(PROCESSED), any(Metadata.class));
verify(retriableStreamRecorder, never()).newSubstream(anyInt());
}

@Test
public void retry_everythingDrained() {
ClientStream mockStream1 = mock(ClientStream.class);
Expand Down
Loading