• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

grpc / grpc-java / #20500

05 Oct 2026 10:45AM UTC coverage: 89.369% (+0.01%) from 89.356%
#20500

push

github

web-flow
core: prevent RetriableStream registry leak on cancel race (#13083)

When `ClientCall.cancel()` races `ClientCall.start()` on a retry-enabled
channel, `RetriableStream.cancel()` can commit the stream before
`start()` calls `prestart()`.

The commit callback then removes nothing because the stream has not been
registered yet. `prestart()` subsequently registers the
already-committed stream, and the one-shot commit prevents any later
removal. This permanently retains the stream in
`UncommittedRetriableStreamsRegistry` and can delay channel shutdown.

This change makes `RetriableStream.start()` check whether the stream was
committed after `prestart()`. If so, it runs `postCommit()` again after
registration, allowing the registry to remove the stream.

Added regression coverage for retry streams, hedging streams, and
streams without retry or hedging policy.

Related issue: #13034

Testing:
- `git diff --check` passed.
- The targeted local Gradle test could not resolve Gradle Plugin Portal
dependencies because DNS resolution failed in the Docker environment.
- The GitHub Actions JDK 17 job failed in the unrelated
`grpc-servlet-jakarta` test
`JettyTransportTest.frameAfterRstStreamShouldNotBreakClientChannel`; the
other reported checks passed.

39266 of 43937 relevant lines covered (89.37%)

0.89 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

95.53
/../core/src/main/java/io/grpc/internal/RetriableStream.java
1
/*
2
 * Copyright 2017 The gRPC Authors
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 */
16

17
package io.grpc.internal;
18

19
import static com.google.common.base.Preconditions.checkArgument;
20
import static com.google.common.base.Preconditions.checkNotNull;
21
import static com.google.common.base.Preconditions.checkState;
22

23
import com.google.common.annotations.VisibleForTesting;
24
import com.google.common.base.Objects;
25
import com.google.errorprone.annotations.CheckReturnValue;
26
import com.google.errorprone.annotations.concurrent.GuardedBy;
27
import io.grpc.Attributes;
28
import io.grpc.ClientStreamTracer;
29
import io.grpc.Compressor;
30
import io.grpc.Deadline;
31
import io.grpc.DecompressorRegistry;
32
import io.grpc.Metadata;
33
import io.grpc.MethodDescriptor;
34
import io.grpc.Status;
35
import io.grpc.SynchronizationContext;
36
import io.grpc.internal.ClientStreamListener.RpcProgress;
37
import java.io.InputStream;
38
import java.lang.Thread.UncaughtExceptionHandler;
39
import java.util.ArrayList;
40
import java.util.Collection;
41
import java.util.Collections;
42
import java.util.List;
43
import java.util.Random;
44
import java.util.concurrent.Executor;
45
import java.util.concurrent.Future;
46
import java.util.concurrent.ScheduledExecutorService;
47
import java.util.concurrent.TimeUnit;
48
import java.util.concurrent.atomic.AtomicBoolean;
49
import java.util.concurrent.atomic.AtomicInteger;
50
import java.util.concurrent.atomic.AtomicLong;
51
import javax.annotation.CheckForNull;
52
import javax.annotation.Nullable;
53

54
/** A logical {@link ClientStream} that is retriable. */
55
abstract class RetriableStream<ReqT> implements ClientStream {
56
  @VisibleForTesting
57
  static final Metadata.Key<String> GRPC_PREVIOUS_RPC_ATTEMPTS =
1 ✔
58
      Metadata.Key.of("grpc-previous-rpc-attempts", Metadata.ASCII_STRING_MARSHALLER);
1 ✔
59

60
  @VisibleForTesting
61
  static final Metadata.Key<String> GRPC_RETRY_PUSHBACK_MS =
1 ✔
62
      Metadata.Key.of("grpc-retry-pushback-ms", Metadata.ASCII_STRING_MARSHALLER);
1 ✔
63

64
  private static final Status CANCELLED_BECAUSE_COMMITTED =
1 ✔
65
      Status.CANCELLED.withDescription("Stream thrown away because RetriableStream committed");
1 ✔
66

67
  private final MethodDescriptor<ReqT, ?> method;
68
  private final Executor callExecutor;
69
  private final Executor listenerSerializeExecutor = new SynchronizationContext(
1 ✔
70
      new UncaughtExceptionHandler() {
1 ✔
71
        @Override
72
        public void uncaughtException(Thread t, Throwable e) {
73
          throw Status.fromThrowable(e)
1 ✔
74
              .withDescription("Uncaught exception in the SynchronizationContext. Re-thrown.")
1 ✔
75
              .asRuntimeException();
1 ✔
76
        }
77
      }
78
  );
79
  private final ScheduledExecutorService scheduledExecutorService;
80
  // Must not modify it.
81
  private final Metadata headers;
82
  @Nullable
83
  private final RetryPolicy retryPolicy;
84
  @Nullable
85
  private final HedgingPolicy hedgingPolicy;
86
  private final boolean isHedging;
87

88
  /** Must be held when updating state, accessing state.buffer, or certain substream attributes. */
89
  private final Object lock = new Object();
1 ✔
90

91
  private final ChannelBufferMeter channelBufferUsed;
92
  private final long perRpcBufferLimit;
93
  private final long channelBufferLimit;
94
  @Nullable
95
  private final Throttle throttle;
96
  @GuardedBy("lock")
1 ✔
97
  private final InsightBuilder closedSubstreamsInsight = new InsightBuilder();
98

99
  private volatile State state = new State(
1 ✔
100
      new ArrayList<BufferEntry>(8), Collections.<Substream>emptyList(), null, null, false, false,
1 ✔
101
      false, 0);
102

103
  /**
104
   * Either non-local transparent retry happened or reached server's application logic.
105
   *
106
   * <p>Note that local-only transparent retries are unlimited.
107
   */
108
  private final AtomicBoolean noMoreTransparentRetry = new AtomicBoolean();
1 ✔
109
  private final AtomicInteger localOnlyTransparentRetries = new AtomicInteger();
1 ✔
110
  private final AtomicInteger inFlightSubStreams = new AtomicInteger();
1 ✔
111
  private SavedCloseMasterListenerReason savedCloseMasterListenerReason;
112

113
  // Used for recording the share of buffer used for the current call out of the channel buffer.
114
  // This field would not be necessary if there is no channel buffer limit.
115
  @GuardedBy("lock")
116
  private long perRpcBufferUsed;
117

118
  private ClientStreamListener masterListener;
119
  @GuardedBy("lock")
120
  private FutureCanceller scheduledRetry;
121
  @GuardedBy("lock")
122
  private FutureCanceller scheduledHedging;
123
  private long nextBackoffIntervalNanos;
124
  private Status cancellationStatus;
125
  private boolean isClosed;
126

127
  RetriableStream(
128
      MethodDescriptor<ReqT, ?> method, Metadata headers,
129
      ChannelBufferMeter channelBufferUsed, long perRpcBufferLimit, long channelBufferLimit,
130
      Executor callExecutor, ScheduledExecutorService scheduledExecutorService,
131
      @Nullable RetryPolicy retryPolicy, @Nullable HedgingPolicy hedgingPolicy,
132
      @Nullable Throttle throttle) {
1 ✔
133
    this.method = method;
1 ✔
134
    this.channelBufferUsed = channelBufferUsed;
1 ✔
135
    this.perRpcBufferLimit = perRpcBufferLimit;
1 ✔
136
    this.channelBufferLimit = channelBufferLimit;
1 ✔
137
    this.callExecutor = callExecutor;
1 ✔
138
    this.scheduledExecutorService = scheduledExecutorService;
1 ✔
139
    this.headers = headers;
1 ✔
140
    this.retryPolicy = retryPolicy;
1 ✔
141
    if (retryPolicy != null) {
1 ✔
142
      this.nextBackoffIntervalNanos = retryPolicy.initialBackoffNanos;
1 ✔
143
    }
144
    this.hedgingPolicy = hedgingPolicy;
1 ✔
145
    checkArgument(
1 ✔
146
        retryPolicy == null || hedgingPolicy == null,
147
        "Should not provide both retryPolicy and hedgingPolicy");
148
    this.isHedging = hedgingPolicy != null;
1 ✔
149
    this.throttle = throttle;
1 ✔
150
  }
1 ✔
151

152
  @SuppressWarnings("GuardedBy")  // TODO(b/145386688) this.lock==ScheduledCancellor.lock so ok
153
  @Nullable // null if already committed
154
  @CheckReturnValue
155
  private Runnable commit(final Substream winningSubstream) {
156
    synchronized (lock) {
1 ✔
157
      if (state.winningSubstream != null) {
1 ✔
158
        return null;
1 ✔
159
      }
160
      final Collection<Substream> savedDrainedSubstreams = state.drainedSubstreams;
1 ✔
161

162
      state = state.committed(winningSubstream);
1 ✔
163

164
      // subtract the share of this RPC from channelBufferUsed.
165
      channelBufferUsed.addAndGet(-perRpcBufferUsed);
1 ✔
166

167
      final boolean wasCancelled = (scheduledRetry != null) ? scheduledRetry.isCancelled() : false;
1 ✔
168
      final Future<?> retryFuture;
169
      final boolean retryWasScheduled = scheduledRetry != null;
1 ✔
170
      if (retryWasScheduled) {
1 ✔
171
        retryFuture = scheduledRetry.markCancelled();
1 ✔
172
        scheduledRetry = null;
1 ✔
173
      } else {
174
        retryFuture = null;
1 ✔
175
      }
176
      // cancel the scheduled hedging if it is scheduled prior to the commitment
177
      final Future<?> hedgingFuture;
178
      if (scheduledHedging != null) {
1 ✔
179
        hedgingFuture = scheduledHedging.markCancelled();
1 ✔
180
        scheduledHedging = null;
1 ✔
181
      } else {
182
        hedgingFuture = null;
1 ✔
183
      }
184

185
      class CommitTask implements Runnable {
1 ✔
186
        @Override
187
        public void run() {
188
          // For hedging only, not needed for normal retry
189
          for (Substream substream : savedDrainedSubstreams) {
1 ✔
190
            if (substream != winningSubstream) {
1 ✔
191
              substream.stream.cancel(CANCELLED_BECAUSE_COMMITTED);
1 ✔
192
            }
193
          }
1 ✔
194
          if (retryWasScheduled) {
1 ✔
195
            if (retryFuture != null) {
1 ✔
196
              retryFuture.cancel(false);
1 ✔
197
            }
198
            if (!wasCancelled && inFlightSubStreams.decrementAndGet() == Integer.MIN_VALUE) {
1 ✔
199
              assert savedCloseMasterListenerReason != null;
×
200
              listenerSerializeExecutor.execute(
×
201
                  new Runnable() {
×
202
                    @Override
203
                    public void run() {
204
                      isClosed = true;
×
205
                      masterListener.closed(savedCloseMasterListenerReason.status,
×
206
                          savedCloseMasterListenerReason.progress,
×
207
                          savedCloseMasterListenerReason.metadata);
×
208
                    }
×
209
                  });
210
            }
211
          }
212

213
          if (hedgingFuture != null) {
1 ✔
214
            hedgingFuture.cancel(false);
1 ✔
215
          }
216

217
          postCommit();
1 ✔
218
        }
1 ✔
219
      }
220

221
      return new CommitTask();
1 ✔
222
    }
223
  }
224

225
  abstract void postCommit();
226

227
  /**
228
   * Calls commit() and if successful runs the post commit task. Post commit task will be non-null
229
   * for only once. The post commit task cancels other non-winning streams on separate transport
230
   * threads, thus it must be run on the callExecutor to prevent deadlocks between multiple stream
231
   * transports.(issues/10314)
232
   * This method should be called only in subListener callbacks. This guarantees callExecutor
233
   * schedules tasks before master listener closes, which is protected by the inFlightSubStreams
234
   * decorative. That is because:
235
   * For a successful winning stream, other streams won't attempt to close master listener.
236
   * For a cancelled winning stream (noop), other stream won't attempt to close master listener.
237
   * For a failed/closed winning stream, the last closed stream closes the master listener, and
238
   * callExecutor scheduling happens-before that.
239
   */
240
  private void commitAndRun(Substream winningSubstream) {
241
    Runnable postCommitTask = commit(winningSubstream);
1 ✔
242

243
    if (postCommitTask != null) {
1 ✔
244
      callExecutor.execute(postCommitTask);
1 ✔
245
    }
246
  }
1 ✔
247

248
  // returns null means we should not create new sub streams, e.g. cancelled or
249
  // other close condition is met for retriableStream.
250
  @Nullable
251
  private Substream createSubstream(int previousAttemptCount, boolean isTransparentRetry,
252
                                    boolean isHedgedStream) {
253
    int inFlight;
254
    do {
255
      inFlight = inFlightSubStreams.get();
1 ✔
256
      if (inFlight < 0) {
1 ✔
257
        return null;
×
258
      }
259
    } while (!inFlightSubStreams.compareAndSet(inFlight, inFlight + 1));
1 ✔
260
    Substream sub = new Substream(previousAttemptCount);
1 ✔
261
    // one tracer per substream
262
    final ClientStreamTracer bufferSizeTracer = new BufferSizeTracer(sub);
1 ✔
263
    ClientStreamTracer.Factory tracerFactory = new ClientStreamTracer.Factory() {
1 ✔
264
      @Override
265
      public ClientStreamTracer newClientStreamTracer(
266
          ClientStreamTracer.StreamInfo info, Metadata headers) {
267
        return bufferSizeTracer;
1 ✔
268
      }
269
    };
270

271
    Metadata newHeaders = updateHeaders(headers, previousAttemptCount);
1 ✔
272
    // NOTICE: This set _must_ be done before stream.start() and it actually is.
273
    sub.stream = newSubstream(newHeaders, tracerFactory, previousAttemptCount, isTransparentRetry,
1 ✔
274
        isHedgedStream);
275
    return sub;
1 ✔
276
  }
277

278
  /**
279
   * Creates a new physical ClientStream that represents a retry/hedging attempt. The returned
280
   * Client stream is not yet started.
281
   */
282
  abstract ClientStream newSubstream(
283
      Metadata headers, ClientStreamTracer.Factory tracerFactory, int previousAttempts,
284
      boolean isTransparentRetry, boolean isHedgedStream);
285

286
  /** Adds grpc-previous-rpc-attempts in the headers of a retry/hedging RPC. */
287
  @VisibleForTesting
288
  final Metadata updateHeaders(
289
      Metadata originalHeaders, int previousAttemptCount) {
290
    Metadata newHeaders = new Metadata();
1 ✔
291
    newHeaders.merge(originalHeaders);
1 ✔
292
    if (previousAttemptCount > 0) {
1 ✔
293
      newHeaders.put(GRPC_PREVIOUS_RPC_ATTEMPTS, String.valueOf(previousAttemptCount));
1 ✔
294
    }
295
    return newHeaders;
1 ✔
296
  }
297

298
  private void drain(Substream substream) {
299
    int index = 0;
1 ✔
300
    int chunk = 0x80;
1 ✔
301
    List<BufferEntry> list = null;
1 ✔
302
    boolean streamStarted = false;
1 ✔
303
    Runnable onReadyRunnable = null;
1 ✔
304

305
    while (true) {
306
      State savedState;
307

308
      synchronized (lock) {
1 ✔
309
        savedState = state;
1 ✔
310
        if (savedState.winningSubstream != null && savedState.winningSubstream != substream) {
1 ✔
311
          // committed but not me, to be cancelled
312
          break;
1 ✔
313
        }
314
        if (savedState.cancelled) {
1 ✔
315
          break;
1 ✔
316
        }
317
        if (index == savedState.buffer.size()) { // I'm drained
1 ✔
318
          state = savedState.substreamDrained(substream);
1 ✔
319
          if (!isReady()) {
1 ✔
320
            return;
1 ✔
321
          }
322
          onReadyRunnable = new Runnable() {
1 ✔
323
            @Override
324
            public void run() {
325
              if (!isClosed) {
1 ✔
326
                masterListener.onReady();
1 ✔
327
              }
328
            }
1 ✔
329
          };
330
          break;
1 ✔
331
        }
332

333
        if (substream.closed) {
1 ✔
334
          return;
×
335
        }
336

337
        int stop = Math.min(index + chunk, savedState.buffer.size());
1 ✔
338
        if (list == null) {
1 ✔
339
          list = new ArrayList<>(savedState.buffer.subList(index, stop));
1 ✔
340
        } else {
341
          list.clear();
1 ✔
342
          list.addAll(savedState.buffer.subList(index, stop));
1 ✔
343
        }
344
        index = stop;
1 ✔
345
      }
1 ✔
346

347
      for (BufferEntry bufferEntry : list) {
1 ✔
348
        bufferEntry.runWith(substream);
1 ✔
349
        if (bufferEntry instanceof RetriableStream.StartEntry) {
1 ✔
350
          streamStarted = true;
1 ✔
351
        }
352
        savedState = state;
1 ✔
353
        if (savedState.winningSubstream != null && savedState.winningSubstream != substream) {
1 ✔
354
          // committed but not me, to be cancelled
355
          break;
1 ✔
356
        }
357
        if (savedState.cancelled) {
1 ✔
358
          break;
1 ✔
359
        }
360
      }
1 ✔
361
    }
1 ✔
362

363
    if (onReadyRunnable != null) {
1 ✔
364
      listenerSerializeExecutor.execute(onReadyRunnable);
1 ✔
365
      return;
1 ✔
366
    }
367

368
    if (!streamStarted) {
1 ✔
369
      // Start stream so inFlightSubStreams is decremented in Sublistener.closed()
370
      substream.stream.start(new Sublistener(substream));
1 ✔
371
    }
372
    substream.stream.cancel(
1 ✔
373
        state.winningSubstream == substream ? cancellationStatus : CANCELLED_BECAUSE_COMMITTED);
1 ✔
374
  }
1 ✔
375

376
  /**
377
   * Runs pre-start tasks. Returns the Status of shutdown if the channel is shutdown.
378
   */
379
  @CheckReturnValue
380
  @Nullable
381
  abstract Status prestart();
382

383
  class StartEntry implements BufferEntry {
1 ✔
384
    @Override
385
    public void runWith(Substream substream) {
386
      substream.stream.start(new Sublistener(substream));
1 ✔
387
    }
1 ✔
388
  }
389

390
  /** Starts the first RPC attempt. */
391
  @Override
392
  public final void start(ClientStreamListener listener) {
393
    masterListener = listener;
1 ✔
394

395
    Status shutdownStatus = prestart();
1 ✔
396

397
    if (shutdownStatus != null) {
1 ✔
398
      cancel(shutdownStatus);
1 ✔
399
      return;
1 ✔
400
    }
401

402
    // cancel() may have committed this stream before prestart() registered it. In that case the
403
    // post-commit callback ran before registration and could not remove the stream. Run it again
404
    // after registration so the channel's uncommitted stream registry cannot retain this stream.
405
    boolean alreadyCommitted;
406
    synchronized (lock) {
1 ✔
407
      alreadyCommitted = state.winningSubstream != null;
1 ✔
408
      if (!alreadyCommitted) {
1 ✔
409
        state.buffer.add(new StartEntry());
1 ✔
410
      }
411
    }
1 ✔
412
    if (alreadyCommitted) {
1 ✔
413
      postCommit();
1 ✔
414
      return;
1 ✔
415
    }
416

417
    Substream substream = createSubstream(0, false, false);
1 ✔
418
    if (substream == null) {
1 ✔
419
      return;
×
420
    }
421
    if (isHedging) {
1 ✔
422
      FutureCanceller scheduledHedgingRef = null;
1 ✔
423

424
      synchronized (lock) {
1 ✔
425
        state = state.addActiveHedge(substream);
1 ✔
426
        if (hasPotentialHedging(state)
1 ✔
427
            && (throttle == null || throttle.isAboveThreshold())) {
1 ✔
428
          scheduledHedging = scheduledHedgingRef = new FutureCanceller(lock);
1 ✔
429
        }
430
      }
1 ✔
431

432
      if (scheduledHedgingRef != null) {
1 ✔
433
        scheduledHedgingRef.setFuture(
1 ✔
434
            scheduledExecutorService.schedule(
1 ✔
435
                new HedgingRunnable(scheduledHedgingRef),
436
                hedgingPolicy.hedgingDelayNanos,
437
                TimeUnit.NANOSECONDS));
438
      }
439
    }
440

441
    drain(substream);
1 ✔
442
  }
1 ✔
443

444
  @SuppressWarnings("GuardedBy")  // TODO(b/145386688) this.lock==ScheduledCancellor.lock so ok
445
  private void pushbackHedging(@Nullable Integer delayMillis) {
446
    if (delayMillis == null) {
1 ✔
447
      return;
1 ✔
448
    }
449
    if (delayMillis < 0) {
1 ✔
450
      freezeHedging();
1 ✔
451
      return;
1 ✔
452
    }
453

454
    // Cancels the current scheduledHedging and reschedules a new one.
455
    FutureCanceller future;
456
    Future<?> futureToBeCancelled;
457

458
    synchronized (lock) {
1 ✔
459
      if (scheduledHedging == null) {
1 ✔
460
        return;
×
461
      }
462

463
      futureToBeCancelled = scheduledHedging.markCancelled();
1 ✔
464
      scheduledHedging = future = new FutureCanceller(lock);
1 ✔
465
    }
1 ✔
466

467
    if (futureToBeCancelled != null) {
1 ✔
468
      futureToBeCancelled.cancel(false);
1 ✔
469
    }
470
    future.setFuture(scheduledExecutorService.schedule(
1 ✔
471
        new HedgingRunnable(future), delayMillis, TimeUnit.MILLISECONDS));
1 ✔
472
  }
1 ✔
473

474
  private final class HedgingRunnable implements Runnable {
475

476
    // Need to hold a ref to the FutureCanceller in case RetriableStrea.scheduledHedging is renewed
477
    // by a positive push-back just after newSubstream is instantiated, so that we can double check.
478
    final FutureCanceller scheduledHedgingRef;
479

480
    HedgingRunnable(FutureCanceller scheduledHedging) {
1 ✔
481
      scheduledHedgingRef = scheduledHedging;
1 ✔
482
    }
1 ✔
483

484
    @Override
485
    public void run() {
486
      // It's safe to read state.hedgingAttemptCount here.
487
      // If this run is not cancelled, the value of state.hedgingAttemptCount won't change
488
      // until state.addActiveHedge() is called subsequently, even the state could possibly
489
      // change.
490
      Substream newSubstream = createSubstream(state.hedgingAttemptCount, false, true);
1 ✔
491
      if (newSubstream == null) {
1 ✔
492
        return;
×
493
      }
494
      callExecutor.execute(
1 ✔
495
          new Runnable() {
1 ✔
496
            @SuppressWarnings("GuardedBy")  //TODO(b/145386688) lock==ScheduledCancellor.lock so ok
497
            @Override
498
            public void run() {
499
              boolean cancelled = false;
1 ✔
500
              FutureCanceller future = null;
1 ✔
501

502
              synchronized (lock) {
1 ✔
503
                if (scheduledHedgingRef.isCancelled()) {
1 ✔
504
                  cancelled = true;
×
505
                } else {
506
                  state = state.addActiveHedge(newSubstream);
1 ✔
507
                  if (hasPotentialHedging(state)
1 ✔
508
                      && (throttle == null || throttle.isAboveThreshold())) {
1 ✔
509
                    scheduledHedging = future = new FutureCanceller(lock);
1 ✔
510
                  } else {
511
                    state = state.freezeHedging();
1 ✔
512
                    scheduledHedging = null;
1 ✔
513
                  }
514
                }
515
              }
1 ✔
516

517
              if (cancelled) {
1 ✔
518
                // Start stream so inFlightSubStreams is decremented in Sublistener.closed()
519
                newSubstream.stream.start(new Sublistener(newSubstream));
×
520
                newSubstream.stream.cancel(Status.CANCELLED.withDescription("Unneeded hedging"));
×
521
                return;
×
522
              }
523
              if (future != null) {
1 ✔
524
                future.setFuture(
1 ✔
525
                    scheduledExecutorService.schedule(
1 ✔
526
                        new HedgingRunnable(future),
527
                        hedgingPolicy.hedgingDelayNanos,
1 ✔
528
                        TimeUnit.NANOSECONDS));
529
              }
530
              drain(newSubstream);
1 ✔
531
            }
1 ✔
532
          });
533
    }
1 ✔
534
  }
535

536
  @Override
537
  public final void cancel(final Status reason) {
538
    Substream noopSubstream = new Substream(0 /* previousAttempts doesn't matter here */);
1 ✔
539
    noopSubstream.stream = new NoopClientStream();
1 ✔
540
    Runnable runnable = commit(noopSubstream);
1 ✔
541

542
    if (runnable != null) {
1 ✔
543
      synchronized (lock) {
1 ✔
544
        state = state.substreamDrained(noopSubstream);
1 ✔
545
      }
1 ✔
546
      runnable.run();
1 ✔
547
      safeCloseMasterListener(reason, RpcProgress.PROCESSED, new Metadata());
1 ✔
548
      return;
1 ✔
549
    }
550

551
    Substream winningSubstreamToCancel = null;
1 ✔
552
    synchronized (lock) {
1 ✔
553
      if (state.drainedSubstreams.contains(state.winningSubstream)) {
1 ✔
554
        winningSubstreamToCancel = state.winningSubstream;
1 ✔
555
      } else { // the winningSubstream will be cancelled while draining
556
        cancellationStatus = reason;
1 ✔
557
      }
558
      state = state.cancelled();
1 ✔
559
    }
1 ✔
560
    if (winningSubstreamToCancel != null) {
1 ✔
561
      winningSubstreamToCancel.stream.cancel(reason);
1 ✔
562
    }
563
  }
1 ✔
564

565
  private void delayOrExecute(BufferEntry bufferEntry) {
566
    Collection<Substream> savedDrainedSubstreams;
567
    synchronized (lock) {
1 ✔
568
      if (!state.passThrough) {
1 ✔
569
        state.buffer.add(bufferEntry);
1 ✔
570
      }
571
      savedDrainedSubstreams = state.drainedSubstreams;
1 ✔
572
    }
1 ✔
573

574
    for (Substream substream : savedDrainedSubstreams) {
1 ✔
575
      bufferEntry.runWith(substream);
1 ✔
576
    }
1 ✔
577
  }
1 ✔
578

579
  /**
580
   * Do not use it directly. Use {@link #sendMessage(Object)} instead because we don't use
581
   * InputStream for buffering.
582
   */
583
  @Override
584
  public final void writeMessage(InputStream message) {
585
    throw new IllegalStateException("RetriableStream.writeMessage() should not be called directly");
×
586
  }
587

588
  final void sendMessage(final ReqT message) {
589
    State savedState = state;
1 ✔
590
    if (savedState.passThrough) {
1 ✔
591
      savedState.winningSubstream.stream.writeMessage(method.streamRequest(message));
1 ✔
592
      return;
1 ✔
593
    }
594

595
    class SendMessageEntry implements BufferEntry {
1 ✔
596
      @Override
597
      public void runWith(Substream substream) {
598
        substream.stream.writeMessage(method.streamRequest(message));
1 ✔
599
        // TODO(ejona): Workaround Netty memory leak. Message writes always need to be followed by
600
        // flushes (or half close), but retry appears to have a code path that the flushes may
601
        // not happen. The code needs to be fixed and this removed. See #9340.
602
        substream.stream.flush();
1 ✔
603
      }
1 ✔
604
    }
605

606
    delayOrExecute(new SendMessageEntry());
1 ✔
607
  }
1 ✔
608

609
  @Override
610
  public final void request(final int numMessages) {
611
    State savedState = state;
1 ✔
612
    if (savedState.passThrough) {
1 ✔
613
      savedState.winningSubstream.stream.request(numMessages);
1 ✔
614
      return;
1 ✔
615
    }
616

617
    class RequestEntry implements BufferEntry {
1 ✔
618
      @Override
619
      public void runWith(Substream substream) {
620
        substream.stream.request(numMessages);
1 ✔
621
      }
1 ✔
622
    }
623

624
    delayOrExecute(new RequestEntry());
1 ✔
625
  }
1 ✔
626

627
  @Override
628
  public final void flush() {
629
    State savedState = state;
1 ✔
630
    if (savedState.passThrough) {
1 ✔
631
      savedState.winningSubstream.stream.flush();
1 ✔
632
      return;
1 ✔
633
    }
634

635
    class FlushEntry implements BufferEntry {
1 ✔
636
      @Override
637
      public void runWith(Substream substream) {
638
        substream.stream.flush();
1 ✔
639
      }
1 ✔
640
    }
641

642
    delayOrExecute(new FlushEntry());
1 ✔
643
  }
1 ✔
644

645
  @Override
646
  public final boolean isReady() {
647
    for (Substream substream : state.drainedSubstreams) {
1 ✔
648
      if (substream.stream.isReady()) {
1 ✔
649
        return true;
1 ✔
650
      }
651
    }
1 ✔
652
    return false;
1 ✔
653
  }
654

655
  @Override
656
  public void optimizeForDirectExecutor() {
657
    class OptimizeDirectEntry implements BufferEntry {
1 ✔
658
      @Override
659
      public void runWith(Substream substream) {
660
        substream.stream.optimizeForDirectExecutor();
1 ✔
661
      }
1 ✔
662
    }
663

664
    delayOrExecute(new OptimizeDirectEntry());
1 ✔
665
  }
1 ✔
666

667
  @Override
668
  public final void setCompressor(final Compressor compressor) {
669
    class CompressorEntry implements BufferEntry {
1 ✔
670
      @Override
671
      public void runWith(Substream substream) {
672
        substream.stream.setCompressor(compressor);
1 ✔
673
      }
1 ✔
674
    }
675

676
    delayOrExecute(new CompressorEntry());
1 ✔
677
  }
1 ✔
678

679
  @Override
680
  public final void setFullStreamDecompression(final boolean fullStreamDecompression) {
681
    class FullStreamDecompressionEntry implements BufferEntry {
1 ✔
682
      @Override
683
      public void runWith(Substream substream) {
684
        substream.stream.setFullStreamDecompression(fullStreamDecompression);
1 ✔
685
      }
1 ✔
686
    }
687

688
    delayOrExecute(new FullStreamDecompressionEntry());
1 ✔
689
  }
1 ✔
690

691
  @Override
692
  public final void setMessageCompression(final boolean enable) {
693
    class MessageCompressionEntry implements BufferEntry {
1 ✔
694
      @Override
695
      public void runWith(Substream substream) {
696
        substream.stream.setMessageCompression(enable);
1 ✔
697
      }
1 ✔
698
    }
699

700
    delayOrExecute(new MessageCompressionEntry());
1 ✔
701
  }
1 ✔
702

703
  @Override
704
  public final void halfClose() {
705
    class HalfCloseEntry implements BufferEntry {
1 ✔
706
      @Override
707
      public void runWith(Substream substream) {
708
        substream.stream.halfClose();
1 ✔
709
      }
1 ✔
710
    }
711

712
    delayOrExecute(new HalfCloseEntry());
1 ✔
713
  }
1 ✔
714

715
  @Override
716
  public final void setAuthority(final String authority) {
717
    class AuthorityEntry implements BufferEntry {
1 ✔
718
      @Override
719
      public void runWith(Substream substream) {
720
        substream.stream.setAuthority(authority);
1 ✔
721
      }
1 ✔
722
    }
723

724
    delayOrExecute(new AuthorityEntry());
1 ✔
725
  }
1 ✔
726

727
  @Override
728
  public final void setDecompressorRegistry(final DecompressorRegistry decompressorRegistry) {
729
    class DecompressorRegistryEntry implements BufferEntry {
1 ✔
730
      @Override
731
      public void runWith(Substream substream) {
732
        substream.stream.setDecompressorRegistry(decompressorRegistry);
1 ✔
733
      }
1 ✔
734
    }
735

736
    delayOrExecute(new DecompressorRegistryEntry());
1 ✔
737
  }
1 ✔
738

739
  @Override
740
  public final void setMaxInboundMessageSize(final int maxSize) {
741
    class MaxInboundMessageSizeEntry implements BufferEntry {
1 ✔
742
      @Override
743
      public void runWith(Substream substream) {
744
        substream.stream.setMaxInboundMessageSize(maxSize);
1 ✔
745
      }
1 ✔
746
    }
747

748
    delayOrExecute(new MaxInboundMessageSizeEntry());
1 ✔
749
  }
1 ✔
750

751
  @Override
752
  public final void setMaxOutboundMessageSize(final int maxSize) {
753
    class MaxOutboundMessageSizeEntry implements BufferEntry {
1 ✔
754
      @Override
755
      public void runWith(Substream substream) {
756
        substream.stream.setMaxOutboundMessageSize(maxSize);
1 ✔
757
      }
1 ✔
758
    }
759

760
    delayOrExecute(new MaxOutboundMessageSizeEntry());
1 ✔
761
  }
1 ✔
762

763
  @Override
764
  public final void setDeadline(final Deadline deadline) {
765
    class DeadlineEntry implements BufferEntry {
1 ✔
766
      @Override
767
      public void runWith(Substream substream) {
768
        substream.stream.setDeadline(deadline);
1 ✔
769
      }
1 ✔
770
    }
771

772
    delayOrExecute(new DeadlineEntry());
1 ✔
773
  }
1 ✔
774

775
  @Override
776
  public final Attributes getAttributes() {
777
    if (state.winningSubstream != null) {
1 ✔
778
      return state.winningSubstream.stream.getAttributes();
1 ✔
779
    }
780
    return Attributes.EMPTY;
×
781
  }
782

783
  @Override
784
  public void appendTimeoutInsight(InsightBuilder insight) {
785
    State currentState;
786
    synchronized (lock) {
1 ✔
787
      insight.appendKeyValue("closed", closedSubstreamsInsight);
1 ✔
788
      currentState = state;
1 ✔
789
    }
1 ✔
790
    if (currentState.winningSubstream != null) {
1 ✔
791
      // TODO(zhangkun83): in this case while other drained substreams have been cancelled in favor
792
      // of the winning substream, they may not have received closed() notifications yet, thus they
793
      // may be missing from closedSubstreamsInsight.  This may be a little confusing to the user.
794
      // We need to figure out how to include them.
795
      InsightBuilder substreamInsight = new InsightBuilder();
1 ✔
796
      currentState.winningSubstream.stream.appendTimeoutInsight(substreamInsight);
1 ✔
797
      insight.appendKeyValue("committed", substreamInsight);
1 ✔
798
    } else {
1 ✔
799
      InsightBuilder openSubstreamsInsight = new InsightBuilder();
1 ✔
800
      // drainedSubstreams doesn't include all open substreams.  Those which have just been created
801
      // and are still catching up with buffered requests (in other words, still draining) will not
802
      // show up.  We think this is benign, because the draining should be typically fast, and it'd
803
      // be indistinguishable from the case where those streams are to be created a little late due
804
      // to delays in the timer.
805
      for (Substream sub : currentState.drainedSubstreams) {
1 ✔
806
        InsightBuilder substreamInsight = new InsightBuilder();
1 ✔
807
        sub.stream.appendTimeoutInsight(substreamInsight);
1 ✔
808
        openSubstreamsInsight.append(substreamInsight);
1 ✔
809
      }
1 ✔
810
      insight.appendKeyValue("open", openSubstreamsInsight);
1 ✔
811
    }
812
  }
1 ✔
813

814
  private static Random random = new Random();
1 ✔
815

816
  @VisibleForTesting
817
  static void setRandom(Random random) {
818
    RetriableStream.random = random;
1 ✔
819
  }
1 ✔
820

821
  /**
822
   * Whether there is any potential hedge at the moment. A false return value implies there is
823
   * absolutely no potential hedge. At least one of the hedges will observe a false return value
824
   * when calling this method, unless otherwise the rpc is committed.
825
   */
826
  // only called when isHedging is true
827
  @GuardedBy("lock")
828
  private boolean hasPotentialHedging(State state) {
829
    return state.winningSubstream == null
1 ✔
830
        && state.hedgingAttemptCount < hedgingPolicy.maxAttempts
831
        && !state.hedgingFrozen;
832
  }
833

834
  @SuppressWarnings("GuardedBy")  // TODO(b/145386688) this.lock==ScheduledCancellor.lock so ok
835
  private void freezeHedging() {
836
    Future<?> futureToBeCancelled = null;
1 ✔
837
    synchronized (lock) {
1 ✔
838
      if (scheduledHedging != null) {
1 ✔
839
        futureToBeCancelled = scheduledHedging.markCancelled();
1 ✔
840
        scheduledHedging = null;
1 ✔
841
      }
842
      state = state.freezeHedging();
1 ✔
843
    }
1 ✔
844

845
    if (futureToBeCancelled != null) {
1 ✔
846
      futureToBeCancelled.cancel(false);
1 ✔
847
    }
848
  }
1 ✔
849

850
  private void safeCloseMasterListener(Status status, RpcProgress progress, Metadata metadata) {
851
    savedCloseMasterListenerReason = new SavedCloseMasterListenerReason(status, progress,
1 ✔
852
        metadata);
853
    if (inFlightSubStreams.addAndGet(Integer.MIN_VALUE) == Integer.MIN_VALUE) {
1 ✔
854
      listenerSerializeExecutor.execute(
1 ✔
855
          new Runnable() {
1 ✔
856
            @Override
857
            public void run() {
858
              isClosed = true;
1 ✔
859
              masterListener.closed(status, progress, metadata);
1 ✔
860
            }
1 ✔
861
          });
862
    }
863
  }
1 ✔
864

865
  private static final boolean isExperimentalRetryJitterEnabled = GrpcUtil
1 ✔
866
          .getFlag("GRPC_EXPERIMENTAL_XDS_RLS_LB", true);
1 ✔
867

868
  public static long intervalWithJitter(long intervalNanos) {
869
    double inverseJitterFactor = isExperimentalRetryJitterEnabled
1 ✔
870
            ? 0.4 * random.nextDouble() + 0.8 : random.nextDouble();
1 ✔
871
    return (long) (intervalNanos * inverseJitterFactor);
1 ✔
872
  }
873

874
  private static final class SavedCloseMasterListenerReason {
875
    private final Status status;
876
    private final RpcProgress progress;
877
    private final Metadata metadata;
878

879
    SavedCloseMasterListenerReason(Status status, RpcProgress progress, Metadata metadata) {
1 ✔
880
      this.status = status;
1 ✔
881
      this.progress = progress;
1 ✔
882
      this.metadata = metadata;
1 ✔
883
    }
1 ✔
884
  }
885

886
  private interface BufferEntry {
887
    /** Replays the buffer entry with the given stream. */
888
    void runWith(Substream substream);
889
  }
890

891
  private final class Sublistener implements ClientStreamListener {
1 ✔
892
    final Substream substream;
893

894
    Sublistener(Substream substream) {
1 ✔
895
      this.substream = substream;
1 ✔
896
    }
1 ✔
897

898
    @Override
899
    public void headersRead(final Metadata headers) {
900
      if (substream.previousAttemptCount > 0) {
1 ✔
901
        headers.discardAll(GRPC_PREVIOUS_RPC_ATTEMPTS);
1 ✔
902
        headers.put(GRPC_PREVIOUS_RPC_ATTEMPTS, String.valueOf(substream.previousAttemptCount));
1 ✔
903
      }
904
      commitAndRun(substream);
1 ✔
905
      if (state.winningSubstream == substream) {
1 ✔
906
        if (throttle != null) {
1 ✔
907
          throttle.onSuccess();
1 ✔
908
        }
909
        listenerSerializeExecutor.execute(
1 ✔
910
            new Runnable() {
1 ✔
911
              @Override
912
              public void run() {
913
                masterListener.headersRead(headers);
1 ✔
914
              }
1 ✔
915
            });
916
      }
917
    }
1 ✔
918

919
    @Override
920
    public void closed(
921
        final Status status, final RpcProgress rpcProgress, final Metadata trailers) {
922
      synchronized (lock) {
1 ✔
923
        state = state.substreamClosed(substream);
1 ✔
924
        closedSubstreamsInsight.append(status.getCode());
1 ✔
925
      }
1 ✔
926

927
      if (inFlightSubStreams.decrementAndGet() == Integer.MIN_VALUE) {
1 ✔
928
        assert savedCloseMasterListenerReason != null;
1 ✔
929
        listenerSerializeExecutor.execute(
1 ✔
930
            new Runnable() {
1 ✔
931
              @Override
932
              public void run() {
933
                isClosed = true;
1 ✔
934
                masterListener.closed(savedCloseMasterListenerReason.status,
1 ✔
935
                    savedCloseMasterListenerReason.progress,
1 ✔
936
                    savedCloseMasterListenerReason.metadata);
1 ✔
937
              }
1 ✔
938
            });
939
        return;
1 ✔
940
      }
941

942
      // handle a race between buffer limit exceeded and closed, when setting
943
      // substream.bufferLimitExceeded = true happens before state.substreamClosed(substream).
944
      if (substream.bufferLimitExceeded) {
1 ✔
945
        commitAndRun(substream);
1 ✔
946
        if (state.winningSubstream == substream) {
1 ✔
947
          safeCloseMasterListener(status, rpcProgress, trailers);
1 ✔
948
        }
949
        return;
1 ✔
950
      }
951
      if (rpcProgress == RpcProgress.MISCARRIED
1 ✔
952
          && localOnlyTransparentRetries.incrementAndGet() > 1_000) {
1 ✔
953
        commitAndRun(substream);
1 ✔
954
        if (state.winningSubstream == substream) {
1 ✔
955
          Status tooManyTransparentRetries = GrpcUtil.statusWithDetails(
1 ✔
956
              Status.Code.INTERNAL, "Too many transparent retries. Might be a bug in gRPC", status);
957
          safeCloseMasterListener(tooManyTransparentRetries, rpcProgress, trailers);
1 ✔
958
        }
959
        return;
1 ✔
960
      }
961

962
      if (state.winningSubstream == null) {
1 ✔
963
        if (rpcProgress == RpcProgress.MISCARRIED
1 ✔
964
            || (rpcProgress == RpcProgress.REFUSED
965
                && noMoreTransparentRetry.compareAndSet(false, true))) {
1 ✔
966
          // transparent retry
967
          final Substream newSubstream = createSubstream(substream.previousAttemptCount,
1 ✔
968
              true, false);
969
          if (newSubstream == null) {
1 ✔
970
            return;
×
971
          }
972
          if (isHedging) {
1 ✔
973
            synchronized (lock) {
1 ✔
974
              // Although this operation is not done atomically with
975
              // noMoreTransparentRetry.compareAndSet(false, true), it does not change the size() of
976
              // activeHedges, so neither does it affect the commitment decision of other threads,
977
              // nor do the commitment decision making threads affect itself.
978
              state = state.replaceActiveHedge(substream, newSubstream);
1 ✔
979
            }
1 ✔
980
          }
981

982
          callExecutor.execute(new Runnable() {
1 ✔
983
            @Override
984
            public void run() {
985
              drain(newSubstream);
1 ✔
986
            }
1 ✔
987
          });
988
          return;
1 ✔
989
        } else if (rpcProgress == RpcProgress.DROPPED) {
1 ✔
990
          // For normal retry, nothing need be done here, will just commit.
991
          // For hedging, cancel scheduled hedge that is scheduled prior to the drop
992
          if (isHedging) {
1 ✔
993
            freezeHedging();
×
994
          }
995
        } else {
996
          noMoreTransparentRetry.set(true);
1 ✔
997

998
          if (isHedging) {
1 ✔
999
            HedgingPlan hedgingPlan = makeHedgingDecision(status, trailers);
1 ✔
1000
            if (hedgingPlan.isHedgeable) {
1 ✔
1001
              pushbackHedging(hedgingPlan.hedgingPushbackMillis);
1 ✔
1002
            }
1003
            synchronized (lock) {
1 ✔
1004
              state = state.removeActiveHedge(substream);
1 ✔
1005
              // The invariant is whether or not #(Potential Hedge + active hedges) > 0.
1006
              // Once hasPotentialHedging(state) is false, it will always be false, and then
1007
              // #(state.activeHedges) will be decreasing. This guarantees that even there may be
1008
              // multiple concurrent hedges, one of the hedges will end up committed.
1009
              if (hedgingPlan.isHedgeable) {
1 ✔
1010
                if (hasPotentialHedging(state) || !state.activeHedges.isEmpty()) {
1 ✔
1011
                  return;
1 ✔
1012
                }
1013
                // else, no activeHedges, no new hedges possible, try to commit
1014
              } // else, isHedgeable is false, try to commit
1015
            }
1 ✔
1016
          } else {
1 ✔
1017
            RetryPlan retryPlan = makeRetryDecision(status, trailers);
1 ✔
1018
            if (retryPlan.shouldRetry) {
1 ✔
1019
              // retry
1020
              Substream newSubstream = createSubstream(substream.previousAttemptCount + 1,
1 ✔
1021
                  false, false);
1022
              if (newSubstream == null) {
1 ✔
1023
                return;
×
1024
              }
1025
              // The check state.winningSubstream == null, checking if is not already committed, is
1026
              // racy, but is still safe b/c the retry will also handle committed/cancellation
1027
              FutureCanceller scheduledRetryCopy;
1028
              synchronized (lock) {
1 ✔
1029
                scheduledRetry = scheduledRetryCopy = new FutureCanceller(lock);
1 ✔
1030
              }
1 ✔
1031

1032
              class RetryBackoffRunnable implements Runnable {
1 ✔
1033
                @Override
1034
                @SuppressWarnings("FutureReturnValueIgnored")
1035
                public void run() {
1036
                  synchronized (scheduledRetryCopy.lock) {
1 ✔
1037
                    if (scheduledRetryCopy.isCancelled()) {
1 ✔
1038
                      return;
×
1039
                    } else {
1040
                      scheduledRetryCopy.markCancelled();
1 ✔
1041
                    }
1042
                  }
1 ✔
1043

1044
                  callExecutor.execute(
1 ✔
1045
                      new Runnable() {
1 ✔
1046
                        @Override
1047
                        public void run() {
1048
                          drain(newSubstream);
1 ✔
1049
                        }
1 ✔
1050
                      });
1051
                }
1 ✔
1052
              }
1053

1054
              scheduledRetryCopy.setFuture(
1 ✔
1055
                  scheduledExecutorService.schedule(
1 ✔
1056
                      new RetryBackoffRunnable(),
1057
                      retryPlan.backoffNanos,
1058
                      TimeUnit.NANOSECONDS));
1059
              return;
1 ✔
1060
            }
1061
          }
1062
        }
1063
      }
1064

1065
      commitAndRun(substream);
1 ✔
1066
      if (state.winningSubstream == substream) {
1 ✔
1067
        safeCloseMasterListener(status, rpcProgress, trailers);
1 ✔
1068
      }
1069
    }
1 ✔
1070

1071
    /**
1072
     * Decides in current situation whether or not the RPC should retry and if it should retry how
1073
     * long the backoff should be. The decision does not take the commitment status into account, so
1074
     * caller should check it separately. It also updates the throttle. It does not change state.
1075
     */
1076
    private RetryPlan makeRetryDecision(Status status, Metadata trailer) {
1077
      if (retryPolicy == null) {
1 ✔
1078
        return new RetryPlan(false, 0);
1 ✔
1079
      }
1080
      boolean shouldRetry = false;
1 ✔
1081
      long backoffNanos = 0L;
1 ✔
1082
      boolean isRetryableStatusCode = retryPolicy.retryableStatusCodes.contains(status.getCode());
1 ✔
1083
      Integer pushbackMillis = getPushbackMills(trailer);
1 ✔
1084
      boolean isThrottled = false;
1 ✔
1085
      if (throttle != null) {
1 ✔
1086
        if (isRetryableStatusCode || (pushbackMillis != null && pushbackMillis < 0)) {
1 ✔
1087
          isThrottled = !throttle.onQualifiedFailureThenCheckIsAboveThreshold();
1 ✔
1088
        }
1089
      }
1090

1091
      if (retryPolicy.maxAttempts > substream.previousAttemptCount + 1 && !isThrottled) {
1 ✔
1092
        if (pushbackMillis == null) {
1 ✔
1093
          if (isRetryableStatusCode) {
1 ✔
1094
            shouldRetry = true;
1 ✔
1095
            backoffNanos = intervalWithJitter(nextBackoffIntervalNanos);
1 ✔
1096
            nextBackoffIntervalNanos = Math.min(
1 ✔
1097
                (long) (nextBackoffIntervalNanos * retryPolicy.backoffMultiplier),
1 ✔
1098
                retryPolicy.maxBackoffNanos);
1 ✔
1099
          } // else no retry
1100
        } else if (pushbackMillis >= 0) {
1 ✔
1101
          shouldRetry = true;
1 ✔
1102
          backoffNanos = TimeUnit.MILLISECONDS.toNanos(pushbackMillis);
1 ✔
1103
          nextBackoffIntervalNanos = retryPolicy.initialBackoffNanos;
1 ✔
1104
        } // else no retry
1105
      } // else no retry
1106

1107
      return new RetryPlan(shouldRetry, backoffNanos);
1 ✔
1108
    }
1109

1110
    private HedgingPlan makeHedgingDecision(Status status, Metadata trailer) {
1111
      Integer pushbackMillis = getPushbackMills(trailer);
1 ✔
1112
      boolean isFatal = !hedgingPolicy.nonFatalStatusCodes.contains(status.getCode());
1 ✔
1113
      boolean isThrottled = false;
1 ✔
1114
      if (throttle != null) {
1 ✔
1115
        if (!isFatal || (pushbackMillis != null && pushbackMillis < 0)) {
1 ✔
1116
          isThrottled = !throttle.onQualifiedFailureThenCheckIsAboveThreshold();
1 ✔
1117
        }
1118
      }
1119
      if (!isFatal && !isThrottled && !status.isOk()
1 ✔
1120
          && (pushbackMillis != null && pushbackMillis > 0)) {
1 ✔
1121
        pushbackMillis = 0; // We want the retry after a nonfatal error to be immediate
1 ✔
1122
      }
1123
      return new HedgingPlan(!isFatal && !isThrottled, pushbackMillis);
1 ✔
1124
    }
1125

1126
    @Nullable
1127
    private Integer getPushbackMills(Metadata trailer) {
1128
      String pushbackStr = trailer.get(GRPC_RETRY_PUSHBACK_MS);
1 ✔
1129
      Integer pushbackMillis = null;
1 ✔
1130
      if (pushbackStr != null) {
1 ✔
1131
        try {
1132
          pushbackMillis = Integer.valueOf(pushbackStr);
1 ✔
1133
        } catch (NumberFormatException e) {
1 ✔
1134
          pushbackMillis = -1;
1 ✔
1135
        }
1 ✔
1136
      }
1137
      return pushbackMillis;
1 ✔
1138
    }
1139

1140
    @Override
1141
    public void messagesAvailable(final MessageProducer producer) {
1142
      State savedState = state;
1 ✔
1143
      checkState(
1 ✔
1144
          savedState.winningSubstream != null, "Headers should be received prior to messages.");
1145
      if (savedState.winningSubstream != substream) {
1 ✔
1146
        GrpcUtil.closeQuietly(producer);
1 ✔
1147
        return;
1 ✔
1148
      }
1149
      listenerSerializeExecutor.execute(
1 ✔
1150
          new Runnable() {
1 ✔
1151
            @Override
1152
            public void run() {
1153
              masterListener.messagesAvailable(producer);
1 ✔
1154
            }
1 ✔
1155
          });
1156
    }
1 ✔
1157

1158
    @Override
1159
    public void onReady() {
1160
      // FIXME(#7089): hedging case is broken.
1161
      if (!isReady()) {
1 ✔
1162
        return;
1 ✔
1163
      }
1164
      listenerSerializeExecutor.execute(
1 ✔
1165
          new Runnable() {
1 ✔
1166
            @Override
1167
            public void run() {
1168
              if (!isClosed) {
1 ✔
1169
                masterListener.onReady();
1 ✔
1170
              }
1171
            }
1 ✔
1172
          });
1173
    }
1 ✔
1174
  }
1175

1176
  private static final class State {
1177
    /** Committed and the winning substream drained. */
1178
    final boolean passThrough;
1179

1180
    /** A list of buffered ClientStream runnables. Set to Null once passThrough. */
1181
    @Nullable final List<BufferEntry> buffer;
1182

1183
    /**
1184
     * Unmodifiable collection of all the open substreams that are drained. Singleton once
1185
     * passThrough; Empty if committed but not passTrough.
1186
     */
1187
    final Collection<Substream> drainedSubstreams;
1188

1189
    /**
1190
     * Unmodifiable collection of all the active hedging substreams.
1191
     *
1192
     * <p>A substream even with the attribute substream.closed being true may be considered still
1193
     * "active" at the moment as long as it is in this collection.
1194
     */
1195
    final Collection<Substream> activeHedges; // not null once isHedging = true
1196

1197
    final int hedgingAttemptCount;
1198

1199
    /** Null until committed. */
1200
    @Nullable final Substream winningSubstream;
1201

1202
    /** Not required to set to true when cancelled, but can short-circuit the draining process. */
1203
    final boolean cancelled;
1204

1205
    /** No more hedging due to events like drop or pushback. */
1206
    final boolean hedgingFrozen;
1207

1208
    State(
1209
        @Nullable List<BufferEntry> buffer,
1210
        Collection<Substream> drainedSubstreams,
1211
        Collection<Substream> activeHedges,
1212
        @Nullable Substream winningSubstream,
1213
        boolean cancelled,
1214
        boolean passThrough,
1215
        boolean hedgingFrozen,
1216
        int hedgingAttemptCount) {
1 ✔
1217
      this.buffer = buffer;
1 ✔
1218
      this.drainedSubstreams =
1 ✔
1219
          checkNotNull(drainedSubstreams, "drainedSubstreams");
1 ✔
1220
      this.winningSubstream = winningSubstream;
1 ✔
1221
      this.activeHedges = activeHedges;
1 ✔
1222
      this.cancelled = cancelled;
1 ✔
1223
      this.passThrough = passThrough;
1 ✔
1224
      this.hedgingFrozen = hedgingFrozen;
1 ✔
1225
      this.hedgingAttemptCount = hedgingAttemptCount;
1 ✔
1226

1227
      checkState(!passThrough || buffer == null, "passThrough should imply buffer is null");
1 ✔
1228
      checkState(
1 ✔
1229
          !passThrough || winningSubstream != null,
1230
          "passThrough should imply winningSubstream != null");
1231
      checkState(
1 ✔
1232
          !passThrough
1233
              || (drainedSubstreams.size() == 1 && drainedSubstreams.contains(winningSubstream))
1 ✔
1234
              || (drainedSubstreams.size() == 0 && winningSubstream.closed),
1 ✔
1235
          "passThrough should imply winningSubstream is drained");
1236
      checkState(!cancelled || winningSubstream != null, "cancelled should imply committed");
1 ✔
1237
    }
1 ✔
1238

1239
    @CheckReturnValue
1240
    // GuardedBy RetriableStream.lock
1241
    State cancelled() {
1242
      return new State(
1 ✔
1243
          buffer, drainedSubstreams, activeHedges, winningSubstream, true, passThrough,
1244
          hedgingFrozen, hedgingAttemptCount);
1245
    }
1246

1247
    /** The given substream is drained. */
1248
    @CheckReturnValue
1249
    // GuardedBy RetriableStream.lock
1250
    State substreamDrained(Substream substream) {
1251
      checkState(!passThrough, "Already passThrough");
1 ✔
1252

1253
      Collection<Substream> drainedSubstreams;
1254
      
1255
      if (substream.closed) {
1 ✔
1256
        drainedSubstreams = this.drainedSubstreams;
1 ✔
1257
      } else if (this.drainedSubstreams.isEmpty()) {
1 ✔
1258
        // optimize for 0-retry, which is most of the cases.
1259
        drainedSubstreams = Collections.singletonList(substream);
1 ✔
1260
      } else {
1261
        drainedSubstreams = new ArrayList<>(this.drainedSubstreams);
1 ✔
1262
        drainedSubstreams.add(substream);
1 ✔
1263
        drainedSubstreams = Collections.unmodifiableCollection(drainedSubstreams);
1 ✔
1264
      }
1265

1266
      boolean passThrough = winningSubstream != null;
1 ✔
1267

1268
      List<BufferEntry> buffer = this.buffer;
1 ✔
1269
      if (passThrough) {
1 ✔
1270
        checkState(
1 ✔
1271
            winningSubstream == substream, "Another RPC attempt has already committed");
1272
        buffer = null;
1 ✔
1273
      }
1274

1275
      return new State(
1 ✔
1276
          buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1277
          hedgingFrozen, hedgingAttemptCount);
1278
    }
1279

1280
    /** The given substream is closed. */
1281
    @CheckReturnValue
1282
    // GuardedBy RetriableStream.lock
1283
    State substreamClosed(Substream substream) {
1284
      substream.closed = true;
1 ✔
1285
      if (this.drainedSubstreams.contains(substream)) {
1 ✔
1286
        Collection<Substream> drainedSubstreams = new ArrayList<>(this.drainedSubstreams);
1 ✔
1287
        drainedSubstreams.remove(substream);
1 ✔
1288
        drainedSubstreams = Collections.unmodifiableCollection(drainedSubstreams);
1 ✔
1289
        return new State(
1 ✔
1290
            buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1291
            hedgingFrozen, hedgingAttemptCount);
1292
      } else {
1293
        return this;
1 ✔
1294
      }
1295
    }
1296

1297
    @CheckReturnValue
1298
    // GuardedBy RetriableStream.lock
1299
    State committed(Substream winningSubstream) {
1300
      checkState(this.winningSubstream == null, "Already committed");
1 ✔
1301

1302
      boolean passThrough = false;
1 ✔
1303
      List<BufferEntry> buffer = this.buffer;
1 ✔
1304
      Collection<Substream> drainedSubstreams;
1305

1306
      if (this.drainedSubstreams.contains(winningSubstream)) {
1 ✔
1307
        passThrough = true;
1 ✔
1308
        buffer = null;
1 ✔
1309
        drainedSubstreams = Collections.singleton(winningSubstream);
1 ✔
1310
      } else {
1311
        drainedSubstreams = Collections.emptyList();
1 ✔
1312
      }
1313

1314
      return new State(
1 ✔
1315
          buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1316
          hedgingFrozen, hedgingAttemptCount);
1317
    }
1318

1319
    @CheckReturnValue
1320
    // GuardedBy RetriableStream.lock
1321
    State freezeHedging() {
1322
      if (hedgingFrozen) {
1 ✔
1323
        return this;
×
1324
      }
1325
      return new State(
1 ✔
1326
          buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1327
          true, hedgingAttemptCount);
1328
    }
1329

1330
    @CheckReturnValue
1331
    // GuardedBy RetriableStream.lock
1332
    // state.hedgingAttemptCount is modified only here.
1333
    // The method is only called in RetriableStream.start() and HedgingRunnable.run()
1334
    State addActiveHedge(Substream substream) {
1335
      // hasPotentialHedging must be true
1336
      checkState(!hedgingFrozen, "hedging frozen");
1 ✔
1337
      checkState(winningSubstream == null, "already committed");
1 ✔
1338

1339
      Collection<Substream> activeHedges;
1340
      if (this.activeHedges == null) {
1 ✔
1341
        activeHedges = Collections.singleton(substream);
1 ✔
1342
      } else {
1343
        activeHedges = new ArrayList<>(this.activeHedges);
1 ✔
1344
        activeHedges.add(substream);
1 ✔
1345
        activeHedges = Collections.unmodifiableCollection(activeHedges);
1 ✔
1346
      }
1347

1348
      int hedgingAttemptCount = this.hedgingAttemptCount + 1;
1 ✔
1349
      return new State(
1 ✔
1350
          buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1351
          hedgingFrozen, hedgingAttemptCount);
1352
    }
1353

1354
    @CheckReturnValue
1355
    // GuardedBy RetriableStream.lock
1356
    // The method is only called in Sublistener.closed()
1357
    State removeActiveHedge(Substream substream) {
1358
      Collection<Substream> activeHedges = new ArrayList<>(this.activeHedges);
1 ✔
1359
      activeHedges.remove(substream);
1 ✔
1360
      activeHedges = Collections.unmodifiableCollection(activeHedges);
1 ✔
1361

1362
      return new State(
1 ✔
1363
          buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1364
          hedgingFrozen, hedgingAttemptCount);
1365
    }
1366

1367
    @CheckReturnValue
1368
    // GuardedBy RetriableStream.lock
1369
    // The method is only called for transparent retry.
1370
    State replaceActiveHedge(Substream oldOne, Substream newOne) {
1371
      Collection<Substream> activeHedges = new ArrayList<>(this.activeHedges);
1 ✔
1372
      activeHedges.remove(oldOne);
1 ✔
1373
      activeHedges.add(newOne);
1 ✔
1374
      activeHedges = Collections.unmodifiableCollection(activeHedges);
1 ✔
1375

1376
      return new State(
1 ✔
1377
          buffer, drainedSubstreams, activeHedges, winningSubstream, cancelled, passThrough,
1378
          hedgingFrozen, hedgingAttemptCount);
1379
    }
1380
  }
1381

1382
  /**
1383
   * A wrapper of a physical stream of a retry/hedging attempt, that comes with some useful
1384
   *  attributes.
1385
   */
1386
  private static final class Substream {
1387
    ClientStream stream;
1388

1389
    // GuardedBy RetriableStream.lock
1390
    boolean closed;
1391

1392
    // setting to true must be GuardedBy RetriableStream.lock
1393
    boolean bufferLimitExceeded;
1394

1395
    final int previousAttemptCount;
1396

1397
    Substream(int previousAttemptCount) {
1 ✔
1398
      this.previousAttemptCount = previousAttemptCount;
1 ✔
1399
    }
1 ✔
1400
  }
1401

1402

1403
  /**
1404
   * Traces the buffer used by a substream.
1405
   */
1406
  class BufferSizeTracer extends ClientStreamTracer {
1407
    // Each buffer size tracer is dedicated to one specific substream.
1408
    private final Substream substream;
1409

1410
    @GuardedBy("lock")
1411
    long bufferNeeded;
1412

1413
    BufferSizeTracer(Substream substream) {
1 ✔
1414
      this.substream = substream;
1 ✔
1415
    }
1 ✔
1416

1417
    /**
1418
     * A message is sent to the wire, so its reference would be released if no retry or
1419
     * hedging were involved. So at this point we have to hold the reference of the message longer
1420
     * for retry, and we need to increment {@code substream.bufferNeeded}.
1421
     */
1422
    @Override
1423
    public void outboundWireSize(long bytes) {
1424
      if (state.winningSubstream != null) {
1 ✔
1425
        return;
1 ✔
1426
      }
1427

1428
      Runnable postCommitTask = null;
1 ✔
1429

1430
      // TODO(zdapeng): avoid using the same lock for both in-bound and out-bound.
1431
      synchronized (lock) {
1 ✔
1432
        if (state.winningSubstream != null || substream.closed) {
1 ✔
1433
          return;
×
1434
        }
1435
        bufferNeeded += bytes;
1 ✔
1436
        if (bufferNeeded <= perRpcBufferUsed) {
1 ✔
1437
          return;
1 ✔
1438
        }
1439

1440
        if (bufferNeeded > perRpcBufferLimit) {
1 ✔
1441
          substream.bufferLimitExceeded = true;
1 ✔
1442
        } else {
1443
          // Only update channelBufferUsed when perRpcBufferUsed is not exceeding perRpcBufferLimit.
1444
          long savedChannelBufferUsed =
1 ✔
1445
              channelBufferUsed.addAndGet(bufferNeeded - perRpcBufferUsed);
1 ✔
1446
          perRpcBufferUsed = bufferNeeded;
1 ✔
1447

1448
          if (savedChannelBufferUsed > channelBufferLimit) {
1 ✔
1449
            substream.bufferLimitExceeded = true;
1 ✔
1450
          }
1451
        }
1452

1453
        if (substream.bufferLimitExceeded) {
1 ✔
1454
          postCommitTask = commit(substream);
1 ✔
1455
        }
1456
      }
1 ✔
1457

1458
      if (postCommitTask != null) {
1 ✔
1459
        postCommitTask.run();
1 ✔
1460
      }
1461
    }
1 ✔
1462
  }
1463

1464
  /**
1465
   *  Used to keep track of the total amount of memory used to buffer retryable or hedged RPCs for
1466
   *  the Channel. There should be a single instance of it for each channel.
1467
   */
1468
  static final class ChannelBufferMeter {
1 ✔
1469
    private final AtomicLong bufferUsed = new AtomicLong();
1 ✔
1470

1471
    @VisibleForTesting
1472
    long addAndGet(long newBytesUsed) {
1473
      return bufferUsed.addAndGet(newBytesUsed);
1 ✔
1474
    }
1475
  }
1476

1477
  /**
1478
   * Used for retry throttling.
1479
   */
1480
  static final class Throttle {
1481

1482
    private static final int THREE_DECIMAL_PLACES_SCALE_UP = 1000;
1483

1484
    /**
1485
     * 1000 times the maxTokens field of the retryThrottling policy in service config.
1486
     * The number of tokens starts at maxTokens. The token_count will always be between 0 and
1487
     * maxTokens.
1488
     */
1489
    final int maxTokens;
1490

1491
    /**
1492
     * Half of {@code maxTokens}.
1493
     */
1494
    final int threshold;
1495

1496
    /**
1497
     * 1000 times the tokenRatio field of the retryThrottling policy in service config.
1498
     */
1499
    final int tokenRatio;
1500

1501
    final AtomicInteger tokenCount = new AtomicInteger();
1 ✔
1502

1503
    Throttle(float maxTokens, float tokenRatio) {
1 ✔
1504
      // tokenRatio is up to 3 decimal places
1505
      this.tokenRatio = (int) (tokenRatio * THREE_DECIMAL_PLACES_SCALE_UP);
1 ✔
1506
      this.maxTokens = (int) (maxTokens * THREE_DECIMAL_PLACES_SCALE_UP);
1 ✔
1507
      this.threshold = this.maxTokens / 2;
1 ✔
1508
      tokenCount.set(this.maxTokens);
1 ✔
1509
    }
1 ✔
1510

1511
    @VisibleForTesting
1512
    boolean isAboveThreshold() {
1513
      return tokenCount.get() > threshold;
1 ✔
1514
    }
1515

1516
    /**
1517
     * Counts down the token on qualified failure and checks if it is above the threshold
1518
     * atomically. Qualified failure is a failure with a retryable or non-fatal status code or with
1519
     * a not-to-retry pushback.
1520
     */
1521
    @VisibleForTesting
1522
    boolean onQualifiedFailureThenCheckIsAboveThreshold() {
1523
      while (true) {
1524
        int currentCount = tokenCount.get();
1 ✔
1525
        if (currentCount == 0) {
1 ✔
1526
          return false;
1 ✔
1527
        }
1528
        int decremented = currentCount - (1 * THREE_DECIMAL_PLACES_SCALE_UP);
1 ✔
1529
        boolean updated = tokenCount.compareAndSet(currentCount, Math.max(decremented, 0));
1 ✔
1530
        if (updated) {
1 ✔
1531
          return decremented > threshold;
1 ✔
1532
        }
1533
      }
×
1534
    }
1535

1536
    @VisibleForTesting
1537
    void onSuccess() {
1538
      while (true) {
1539
        int currentCount = tokenCount.get();
1 ✔
1540
        if (currentCount == maxTokens) {
1 ✔
1541
          break;
1 ✔
1542
        }
1543
        int incremented = currentCount + tokenRatio;
1 ✔
1544
        boolean updated = tokenCount.compareAndSet(currentCount, Math.min(incremented, maxTokens));
1 ✔
1545
        if (updated) {
1 ✔
1546
          break;
1 ✔
1547
        }
1548
      }
×
1549
    }
1 ✔
1550

1551
    @Override
1552
    public boolean equals(Object o) {
1553
      if (this == o) {
1 ✔
1554
        return true;
×
1555
      }
1556
      if (!(o instanceof Throttle)) {
1 ✔
1557
        return false;
×
1558
      }
1559
      Throttle that = (Throttle) o;
1 ✔
1560
      return maxTokens == that.maxTokens && tokenRatio == that.tokenRatio;
1 ✔
1561
    }
1562

1563
    @Override
1564
    public int hashCode() {
1565
      return Objects.hashCode(maxTokens, tokenRatio);
×
1566
    }
1567
  }
1568

1569
  private static final class RetryPlan {
1570
    final boolean shouldRetry;
1571
    final long backoffNanos;
1572

1573
    RetryPlan(boolean shouldRetry, long backoffNanos) {
1 ✔
1574
      this.shouldRetry = shouldRetry;
1 ✔
1575
      this.backoffNanos = backoffNanos;
1 ✔
1576
    }
1 ✔
1577
  }
1578

1579
  private static final class HedgingPlan {
1580
    final boolean isHedgeable;
1581
    @Nullable
1582
    final Integer hedgingPushbackMillis;
1583

1584
    public HedgingPlan(
1585
        boolean isHedgeable, @Nullable Integer hedgingPushbackMillis) {
1 ✔
1586
      this.isHedgeable = isHedgeable;
1 ✔
1587
      this.hedgingPushbackMillis = hedgingPushbackMillis;
1 ✔
1588
    }
1 ✔
1589
  }
1590

1591
  /** Allows cancelling a Future without racing with setting the future. */
1592
  private static final class FutureCanceller {
1593

1594
    final Object lock;
1595
    @GuardedBy("lock")
1596
    Future<?> future;
1597
    @GuardedBy("lock")
1598
    boolean cancelled;
1599

1600
    FutureCanceller(Object lock) {
1 ✔
1601
      this.lock = lock;
1 ✔
1602
    }
1 ✔
1603

1604
    void setFuture(Future<?> future) {
1605
      boolean wasCancelled;
1606
      synchronized (lock) {
1 ✔
1607
        wasCancelled = cancelled;
1 ✔
1608
        if (!wasCancelled) {
1 ✔
1609
          this.future = future;
1 ✔
1610
        }
1611
      }
1 ✔
1612
      if (wasCancelled) {
1 ✔
1613
        future.cancel(false);
×
1614
      }
1615
    }
1 ✔
1616

1617
    @GuardedBy("lock")
1618
    @CheckForNull // Must cancel the returned future if not null.
1619
    Future<?> markCancelled() {
1620
      cancelled = true;
1 ✔
1621
      return future;
1 ✔
1622
    }
1623

1624
    @GuardedBy("lock")
1625
    boolean isCancelled() {
1626
      return cancelled;
1 ✔
1627
    }
1628
  }
1629
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc