RetriableStream permanently leaked in UncommittedRetriableStreamsRegistry when cancel() races ClientCall.start()
Nessuno ha ancora preso questa issue.
- Lingua principale
- Java
- Stelle
- 12.1k
- Fork
- 4k
- Merge medio
- 2g 17h
- PR unite (30g)
- 37
Descrizione
### What version of gRPC-Java are you using?
Reproduced on **1.84.0** (latest) and **1.81.0**.
`RetriableStream.java`, `ClientCallImpl.java` and `ClientCalls.java` are **byte-identical between `v1.81.0` and `master`**, and the `ManagedChannelImpl` diff over that range does not touch the registry, `prestart()` or `postCommit()` — so `master` is affected too. All line numbers below are `master`.
### What is your environment?
macOS / JDK 17, `grpc-inprocess` transport in the attached reproducer. Originally found in production on Linux/JDK11 over `grpc-netty`; the defect is transport-independent (it is entirely inside `ClientCallImpl` + `RetriableStream`).
### What did you expect to see?
A `ClientCall.cancel()` should never leave the `RetriableStream` registered in `ManagedChannelImpl$UncommittedRetriableStreamsRegistry`, regardless of when the cancel lands relative to `start()`.
### What did you see instead?
On any channel built with `enableRetry()`, a `cancel()` that lands in the window between `stream = clientStreamProvider.newStream(...)` and `stream.start(...)` inside `ClientCallImpl.startInternal()` leaks the `RetriableStream` **permanently** into that channel-scoped, process-lifetime `HashSet`.
The order of operations inverts:
1. `cancel()` commits the stream and runs `postCommit()` → `registry.remove(this)` — **a no-op, because the stream has not been registered yet**;
2. `start()` then runs `prestart()` → `registry.add(this)`;
3. `commit()` is one-shot, so **no second post-commit task can ever be produced**.
The entry is never removed for the life of the channel. Each leaked entry retains its `Substream`, `Metadata`, `CallOptions` and `Context`, so retained size is substantial. A leaked entry also permanently blocks `delayedTransport.shutdown()` on channel shutdown, since that is gated on the registry reaching `isEmpty()`.
This is **not** hedging-specific — it reproduces with hedging disabled, i.e. on any `enableRetry()` channel.
#### Mechanism
`ClientCallImpl.startInternal()`:
```java
250: stream = clientStreamProvider.newStream(method, callOptions, headers, context);
// ... ~8 setters, reportCallStarted() ...
286: stream.start(new ClientStreamListenerImpl(observer)); // -> RetriableStream.start()
292: cancellationHandler.setUp();
```
`RetriableStream.start()` registers only at step 286:
```java
392: public final void start(ClientStreamListener listener) {
393: masterListener = listener;
395: Status shutdownStatus = prestart(); // ManagedChannelImpl: registry.add(this)
```
`ClientCallImpl.cancelInternal()` is explicitly documented to run in this window, and the `stream != null` guard is already true from L250:
```java
469: // Cancel is called in exception handling cases, so it may be the case that the
470: // stream was never successfully created or start has never been called.
471: if (stream != null) {
482: stream.cancel(status); // non-null since L250 -- but not yet registered
```
`RetriableStream.cancel()` commits **and runs the post-commit task inline**:
```java
526: public final void cancel(final Status reason) {
527: Substream noopSubstream = new Substream(0);
528: noopSubstream.stream = new NoopClientStream();
529: Runnable runnable = commit(noopSubstream); // sets state.winningSubstream
532: if (runnable != null) {
535: runnable.run(); // -> postCommit() -> registry.remove(this) == NO-OP
```
and `commit()` is one-shot:
```java
155: private Runnable commit(final Substream winningSubstream) {
156: synchronized (lock) {
157: if (state.winningSubstream != null) {
158: return null; // no further post-commit task, ever
```
`postCommit()` is the **only** caller of `UncommittedRetriableStreamsRegistry.remove(...)`.
### Steps to reproduce the bug
Single file, stock gRPC jars only, no protobuf and no third-party code (source inlined below). It uses the in-process transport, an `enableRetry()` channel, and a second thread calling `call.cancel()` with randomised nanosecond jitter around `call.start()`; it then reflects on `ManagedChannelImpl.uncommittedRetriableStreamsRegistry.uncommittedRetriableStreams.size()`.
```
javac -cp LeakRepro.java
java -cp .: LeakRepro race # leaks
java -cp .: LeakRepro blocking # control, 0
java -cp .: -Dpolicy=none LeakRepro race # hedging off -- still leaks
```
(The mode is a positional arg and defaults to `blocking`; running it with no args exercises the control path.)
#### Results — 200 calls each, after quiescence + `System.gc()`
| run | version | cancel races start() | retry policy | leaked entries |
|---|---|---|---|---|
| A | 1.84.0 | yes | hedging | **22**, **60**, **50** (three runs) |
| B | 1.81.0 | yes | hedging | **54** |
| C | 1.84.0 | yes | **none** (hedging off) | **9** |
| D *(control)* | 1.84.0 | no | hedging | **0** |
| E *(control)* | 1.81.0 | no | hedging | **0** |
The exact count is timing-dependent and varies run to run; what is stable is **leaked > 0 in every race run and exactly 0 in every control run**, and registry size growing monotonically with call count in A/B/C. Run C is what widens the scope: the leak needs only `enableRetry()`, not hedging.
The controls are not vacuous — the same reflection probe reports a non-zero *peak* registry size during run E (`max registry size ever observed = 1`, tracking the single in-flight call), so it demonstrably can see entries and the final `0` is a real empty registry rather than a broken probe.
Every leaked entry has the same shape, which is exactly what the mechanism predicts:
```
leaked entry shape: total=50 committed(winningSubstream!=null)=50 inFlight=0
cancelled=false passThrough=true hedgingFrozen=false hedgingAttemptCount=0
buffer=null drainedSubstreams=size=1 winningSubstream.stream=NoopClientStream
```
`winningSubstream.stream == NoopClientStream` is constructed **only** in `RetriableStream.cancel()` (L528), so each leaked entry provably committed via `cancel()` and was registered afterwards.
The window is measurable: widening the injected jitter from 4 µs to 400 µs takes the leak from 95 entries down to 1, consistent with a window of tens of microseconds that narrows further once `startInternal` is JIT-compiled.
#### Production impact
Found in a high-QPS Java service on 1.81.0: a heap dump held **499,418** leaked `ManagedChannelImpl$ChannelStreamProvider$1RetryStream` instances, roughly **14 GB**, producing a GC death spiral. Because the retention is channel-scoped, the only recovery is a process restart.
### Suggested fix
Register before the cancellable window, or make de-registration idempotent — any of:
1. **Register in the constructor / `newStream()`** rather than in `start()`, so `add()` happens-before any reachable `cancel()`.
2. **Re-check after `prestart()`**: in `RetriableStream.start()`, if the stream is already committed immediately after `prestart()`, call `postCommit()` again.
3. **Do not cancel an unstarted stream in `ClientCallImpl`** — track "started" and defer the cancel until after `stream.start()`.
(2) is the smallest change and is local to `RetriableStream.start()`:
```java
Status shutdownStatus = prestart();
if (shutdownStatus != null) {
cancel(shutdownStatus);
return;
}
+ // cancel() may have committed this stream before prestart() registered it, in which case
+ // its postCommit() removed nothing -- and commit() is one-shot, so there is no second
+ // chance to de-register.
+ synchronized (lock) {
+ if (state.winningSubstream != null) {
+ postCommit();
+ return;
+ }
+ }
```
A regression test can assert that `uncommittedRetriableStreams` is empty after a `cancel()`-racing-`start()` sequence on an `enableRetry()` channel.
### Related issues checked
I could not find an existing issue for this. Searched `uncommittedRetriableStreams` (0 hits), `UncommittedRetriableStreamsRegistry` (#4073 lock granularity, #7227 refactor, #7362 pending-call drain — none this), and `RetriableStream leak cancel` (#10209 framer/ByteBuf, #9340, #9185, #12891, #3839 — none this). #3537 / #3557 concern a different defect (`ThreadlessExecutor` discarding runnables after `shutdown()`), which I initially suspected here and then ruled out — 360,000 hedged blocking calls over Netty produced 0 dropped runnables and 0 registry growth, with `GRPC_CLIENT_CALL_REJECT_RUNNABLE=true` and a passing positive control for the drop detector.
LeakRepro.java (click to expand)
```java
import io.grpc.*;
import io.grpc.inprocess.InProcessChannelBuilder;
import io.grpc.inprocess.InProcessServerBuilder;
import io.grpc.stub.ClientCalls;
import java.io.ByteArrayInputStream;
import java.io.InputStream;
import java.lang.reflect.Field;
import java.nio.charset.StandardCharsets;
import java.util.*;
import java.util.concurrent.*;
/**
* Standalone reproducer for a permanent RetriableStream leak in
* ManagedChannelImpl$UncommittedRetriableStreamsRegistry.
*
* Bug:
* ClientCallImpl.startInternal() assigns `stream = clientStreamProvider.newStream(...)`
* (~L250) but the stream only registers itself much later, from
* RetriableStream.start() -> prestart() -> registry.add(this), reached at ~L286.
* A ClientCall.cancel() landing in that window reaches RetriableStream.cancel(), which
* commits and runs the post-commit task INLINE -> postCommit() -> registry.remove(this),
* which is a NO-OP because the stream is not registered yet. start() then proceeds and
* registers it. commit() is one-shot ("if (state.winningSubstream != null) return null"),
* so no second post-commit task is ever produced and the entry is pinned for the life of
* the channel.
*
* Requires only an enableRetry() channel -- hedging is NOT required (see -Dpolicy=none).
*
* Prediction:
* mode=race -> registry size grows monotonically and stays > 0 after quiescence
* mode=blocking -> registry size returns to 0 (no cancel races start())
*
* Usage: java LeakRepro [iterations]
* NOTE: the mode is a POSITIONAL arg and defaults to "blocking"; running this class with
* no arguments exercises the control path and does NOT reproduce the leak.
* -Dpolicy=none disables hedging while keeping enableRetry().
*/
public class LeakRepro {
static final MethodDescriptor METHOD =
MethodDescriptor.newBuilder()
.setType(MethodDescriptor.MethodType.UNARY)
.setFullMethodName("leak.Svc/Call")
.setRequestMarshaller(new StringMarshaller())
.setResponseMarshaller(new StringMarshaller())
.build();
/** Trivial marshaller so the repro needs no protobuf dependency. */
static final class StringMarshaller implements MethodDescriptor.Marshaller {
@Override public InputStream stream(String value) {
return new ByteArrayInputStream(value.getBytes(StandardCharsets.UTF_8));
}
@Override public String parse(InputStream stream) {
try {
return new String(stream.readAllBytes(), StandardCharsets.UTF_8);
} catch (Exception e) {
throw new RuntimeException(e);
}
}
}
/**
* A handler that accepts the call and then never responds. The client's deadline is
* what ends every call, which is exactly the production shape (callee slow / hedged
* attempts outliving the caller's deadline).
*/
static ServerServiceDefinition blackHoleService() {
ServerCallHandler handler =
(call, headers) -> {
call.request(2);
return new ServerCall.Listener() {};
};
return ServerServiceDefinition.builder("leak.Svc")
.addMethod(METHOD, handler)
.build();
}
/** Hedging policy equivalent to what a D2/service-config-driven client would receive. */
static Map hedgingServiceConfig(String hedgingDelay, int maxAttempts) {
Map name = new LinkedHashMap<>();
name.put("service", "leak.Svc");
Map hedging = new LinkedHashMap<>();
hedging.put("maxAttempts", (double) maxAttempts);
hedging.put("hedgingDelay", hedgingDelay);
hedging.put("nonFatalStatusCodes", Arrays.asList("UNAVAILABLE", "DEADLINE_EXCEEDED"));
Map methodConfig = new LinkedHashMap<>();
methodConfig.put("name", Collections.singletonList(name));
methodConfig.put("hedgingPolicy", hedging);
Map serviceConfig = new LinkedHashMap<>();
serviceConfig.put("methodConfig", Collections.singletonList(methodConfig));
return serviceConfig;
}
/** Retry policy — an alternative to hedging; also wraps the call in RetriableStream. */
static Map retryServiceConfig(int maxAttempts) {
Map name = new LinkedHashMap<>();
name.put("service", "leak.Svc");
Map retry = new LinkedHashMap<>();
retry.put("maxAttempts", (double) maxAttempts);
retry.put("initialBackoff", "0.1s");
retry.put("maxBackoff", "1s");
retry.put("backoffMultiplier", 2.0D);
retry.put("retryableStatusCodes", Arrays.asList("UNAVAILABLE"));
Map methodConfig = new LinkedHashMap<>();
methodConfig.put("name", Collections.singletonList(name));
methodConfig.put("retryPolicy", retry);
Map serviceConfig = new LinkedHashMap<>();
serviceConfig.put("methodConfig", Collections.singletonList(methodConfig));
return serviceConfig;
}
/** Reflect down to ManagedChannelImpl.uncommittedRetriableStreamsRegistry.uncommittedRetriableStreams. */
static Collection registry(ManagedChannel ch) throws Exception {
Object o = ch;
while (!o.getClass().getName().equals("io.grpc.internal.ManagedChannelImpl")) {
Field delegate = null;
for (Class c = o.getClass(); c != null && delegate == null; c = c.getSuperclass()) {
for (Field f : c.getDeclaredFields()) {
if (ManagedChannel.class.isAssignableFrom(f.getType())) {
delegate = f;
break;
}
}
}
if (delegate == null) {
throw new IllegalStateException("cannot unwrap channel class " + o.getClass().getName());
}
delegate.setAccessible(true);
o = delegate.get(o);
}
Field regField = o.getClass().getDeclaredField("uncommittedRetriableStreamsRegistry");
regField.setAccessible(true);
Object reg = regField.get(o);
Field setField = reg.getClass().getDeclaredField("uncommittedRetriableStreams");
setField.setAccessible(true);
return (Collection) setField.get(reg);
}
/**
* For leaked entries, confirm they match the production heap signature:
* RetriableStream.state.winningSubstream != null (i.e. commit() DID run, so the
* stream is not merely "in flight" -- it is committed-but-never-removed).
*/
static String describeLeaked(Collection reg) throws Exception {
int committed = 0, inFlight = 0, total = 0;
StringBuilder sample = new StringBuilder();
for (Object stream : new ArrayList<>(reg)) {
total++;
Field stateF = null;
for (Class c = stream.getClass(); c != null; c = c.getSuperclass()) {
try {
stateF = c.getDeclaredField("state");
break;
} catch (NoSuchFieldException ignored) { }
}
if (stateF == null) continue;
stateF.setAccessible(true);
Object state = stateF.get(stream);
if (state == null) continue;
Field winF = state.getClass().getDeclaredField("winningSubstream");
winF.setAccessible(true);
Object win = winF.get(state);
if (win != null) committed++; else inFlight++;
if (total <= 3) {
// Compare field-for-field against the production heap dump signature.
sample.append(String.format("%n sample#%d: ", total));
for (String fn : new String[]{"cancelled", "passThrough", "hedgingFrozen",
"hedgingAttemptCount", "buffer", "drainedSubstreams"}) {
Field f = state.getClass().getDeclaredField(fn);
f.setAccessible(true);
Object v = f.get(state);
if (fn.equals("drainedSubstreams") && v instanceof Collection) {
v = "size=" + ((Collection) v).size();
} else if (fn.equals("buffer")) {
v = (v == null ? "null" : "size=" + ((Collection) v).size());
}
sample.append(fn).append('=').append(v).append(" ");
}
if (win != null) {
Field sf = win.getClass().getDeclaredField("stream");
sf.setAccessible(true);
Object s = sf.get(win);
sample.append("winningSubstream.stream=")
.append(s == null ? "null" : s.getClass().getSimpleName());
}
}
}
return String.format("total=%d committed(winningSubstream!=null)=%d inFlight=%d%s",
total, committed, inFlight, sample);
}
public static void main(String[] args) throws Exception {
String mode = args.length > 0 ? args[0] : "blocking";
int iterations = args.length > 1 ? Integer.parseInt(args[1]) : 200;
long deadlineMs = 60;
String serverName = InProcessServerBuilder.generateName();
Server server = InProcessServerBuilder.forName(serverName)
.addService(blackHoleService())
.build()
.start();
// -Dpolicy=hedging (default) | retry | none
// "retry" models a fallback where hedging.nonFatalStatusCodes is left blank.
// "none" models enableRetry() with no policy at all — still wraps in RetriableStream.
String policy = System.getProperty("policy", "hedging");
InProcessChannelBuilder cb = InProcessChannelBuilder.forName(serverName);
if ("off".equals(policy)) {
cb.disableRetry(); // no RetriableStream is ever created
} else {
cb.enableRetry();
}
// -DhedgingDelay varies the hedging delay. The delay gates only when a hedge FIRES;
// the timer is scheduled on every call regardless. The leak race happens at start()
// before any hedge fires, so this knob makes no difference -- and -Dpolicy=none, which
// removes hedging entirely, still leaks.
String hedgingDelay = System.getProperty("hedgingDelay", "0.010s");
if ("hedging".equals(policy)) {
cb.maxHedgedAttempts(3).defaultServiceConfig(hedgingServiceConfig(hedgingDelay, 3));
} else if ("retry".equals(policy)) {
cb.maxRetryAttempts(3).defaultServiceConfig(retryServiceConfig(3));
}
ManagedChannel channel = cb.build();
Collection reg = registry(channel);
System.out.printf("mode=%s iterations=%d deadline=%dms policy=%s hedgingDelay=%s%n",
mode, iterations, deadlineMs, policy, hedgingDelay);
System.out.printf("registry class=%s initial size=%d%n",
reg.getClass().getName(), reg.size());
// POSITIVE CONTROL for the probe itself: if the registry never goes above 0 even
// transiently, then no RetriableStream is being created at all (hedging inactive),
// and a final size of 0 says nothing about cleanup. Sample continuously.
final int[] maxSeen = {0};
Thread sampler = new Thread(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
int s = reg.size();
synchronized (maxSeen) { if (s > maxSeen[0]) maxSeen[0] = s; }
Thread.sleep(0, 200_000);
} catch (InterruptedException e) {
return;
} catch (ConcurrentModificationException ignored) {
// HashSet is not thread-safe; a racing size() read may throw. Ignore and resample.
}
}
});
sampler.setDaemon(true);
sampler.start();
int deadlineExceeded = 0, other = 0;
if (mode.equals("race")) {
// Reproduce the cancel-during-startInternal window:
// ClientCallImpl assigns `stream` (line ~250) well before calling stream.start()
// (line ~286), which is where RetriableStream.prestart() registers the stream.
// A cancel() landing in that window runs postCommit() -> registry.remove(this)
// on a stream that has not been added yet, and the subsequent start() then adds
// it -- permanently, because nothing will ever commit it again.
ExecutorService cancellers = Executors.newFixedThreadPool(8);
java.util.Random rnd = new java.util.Random(42);
for (int i = 1; i <= iterations; i++) {
CallOptions opts = CallOptions.DEFAULT.withDeadlineAfter(deadlineMs, TimeUnit.MILLISECONDS);
final ClientCall call = channel.newCall(METHOD, opts);
final long jitter = rnd.nextInt(4000);
cancellers.submit(() -> {
java.util.concurrent.locks.LockSupport.parkNanos(jitter);
try {
call.cancel("race", null);
} catch (RuntimeException ignored) { }
});
try {
call.start(new ClientCall.Listener() { }, new Metadata());
call.request(2);
call.sendMessage("req");
call.halfClose();
other++;
} catch (RuntimeException e) {
other++;
}
if (i % Math.max(1, iterations / 10) == 0) {
System.out.printf(" after %5d calls: registry size = %d%n", i, reg.size());
}
}
cancellers.shutdown();
cancellers.awaitTermination(30, TimeUnit.SECONDS);
} else {
for (int i = 1; i <= iterations; i++) {
CallOptions opts = CallOptions.DEFAULT.withDeadlineAfter(deadlineMs, TimeUnit.MILLISECONDS);
try {
if (mode.equals("blocking")) {
ClientCalls.blockingUnaryCall(channel, METHOD, opts, "req");
} else {
Future f = ClientCalls.futureUnaryCall(channel.newCall(METHOD, opts), "req");
f.get(deadlineMs + 500, TimeUnit.MILLISECONDS);
}
other++;
} catch (StatusRuntimeException e) {
if (e.getStatus().getCode() == Status.Code.DEADLINE_EXCEEDED) deadlineExceeded++; else other++;
} catch (ExecutionException e) {
Throwable c = e.getCause();
if (c instanceof StatusRuntimeException
&& ((StatusRuntimeException) c).getStatus().getCode() == Status.Code.DEADLINE_EXCEEDED) {
deadlineExceeded++;
} else other++;
} catch (Exception e) {
other++;
}
if (i % Math.max(1, iterations / 10) == 0) {
System.out.printf(" after %4d calls: registry size = %d%n", i, reg.size());
}
}
}
// Give any late async cleanup a chance to run, so a slow-but-eventual removal is not
// misread as a permanent leak.
Thread.sleep(3000);
System.gc();
Thread.sleep(1000);
sampler.interrupt();
System.out.printf("PROBE CONTROL: max registry size ever observed during run = %d%n", maxSeen[0]);
System.out.printf("FINAL mode=%s calls=%d deadlineExceeded=%d other=%d registrySize=%d%n",
mode, iterations, deadlineExceeded, other, reg.size());
System.out.println("leaked entry shape: " + describeLeaked(reg));
channel.shutdownNow();
server.shutdownNow();
}
}
```
Guida per i contributori
Apri la guida per i contributori
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Direzione di ricerca
Inizia da RetriableStream.start() e confronta il flusso di prestart(), commit() e postCommit() prima dell'avvio con ClientCallImpl.startInternal() e cancelInternal(). Usa il LeakRepro.java fornito con un canale enableRetry() e la modalità race, quindi aggiungi o esegui un test di regressione che verifichi che il registro degli elementi non sottoposti a commit sia vuoto dopo le race tra la cancellazione e start().
Scritto dal modello di indicizzazione a partire dal testo della issue.
Valutazione
- Stack tecnologico
- java
- Ambito
- api, backend
- Tipo di issue
- Bug
- Difficoltà
- 3/5
- Tempo stimato
- 1-2 giorni
- Stato di attività
- Attiva
- Chiarezza
- Specificata chiaramente
- Idoneità per principianti
- 68/100