From 764a857d760208e1e27db74948beaa65571555ea Mon Sep 17 00:00:00 2001 From: jdcormie Date: Fri, 7 Aug 2026 11:24:52 -0700 Subject: [PATCH] stub: Clarify flow control javadoc Remove vestigial references to credit-based flow control and the suggestion that a request() for messages directly corresponds to the number of messages a peer could subsequently send before its end of the stream went !isReady(). In fact, okhttp, netty and binder transports all use buffers meaning the sender's stream can be isReady() even when the receiver has no outstanding request()s for messages. Furthermore, those buffers are sized in bytes and messages are not all the same size. So consuming an inbound's next message may not cause a non-ready outbound to become ready again. Replace this with a short discussion of what is actually guaranteed by every transport. --- .../java/io/grpc/stub/CallStreamObserver.java | 27 +++++++++++++++---- .../grpc/stub/ClientCallStreamObserver.java | 3 +-- .../grpc/stub/ServerCallStreamObserver.java | 3 +-- 3 files changed, 24 insertions(+), 9 deletions(-) diff --git a/stub/src/main/java/io/grpc/stub/CallStreamObserver.java b/stub/src/main/java/io/grpc/stub/CallStreamObserver.java index 52dd046831d..3ab9d566f38 100644 --- a/stub/src/main/java/io/grpc/stub/CallStreamObserver.java +++ b/stub/src/main/java/io/grpc/stub/CallStreamObserver.java @@ -47,6 +47,23 @@ *

Like {@code StreamObserver}, implementations are not required to be thread-safe; if multiple * threads will be writing to an instance concurrently, the application must synchronize its calls. * + *

On flow control: The {@link #isReady} state of an outbound {@link CallStreamObserver} is + * related to the number and size of messages written to its {@link StreamObserver#onNext} and + * whether the peer has consumed those messages by way of {@link StreamObserver#onNext} on the + * corresponding inbound. However, this effect may be delayed and is not guaranteed to be reflected + * message-for-message. + * + *

What's actually guaranteed: + * + *

    + *
  1. If an application keeps writing messages to an outbound {@link StreamObserver#onNext} but + * the peer stops consuming them (whether by not returning from the {@link + * StreamObserver#onNext} callback or by not {@link #request}ing those callbacks), then + * eventually the outbound's {@link #isReady} will become false and stay that way. + *
  2. Once in that state of backpressure, requesting and consuming enough messages from the + * 'inbound' end will eventually cause the outbound to become ready again. + *
+ * *

DO NOT MOCK: The API is too complex to reliably mock. Use InProcessChannelBuilder to create * "real" RPCs suitable for testing. * @@ -87,9 +104,10 @@ public abstract class CallStreamObserver implements StreamObserver { public abstract void setOnReadyHandler(Runnable onReadyHandler); /** - * Disables automatic flow control where a token is returned to the peer after a call - * to the 'inbound' {@link io.grpc.stub.StreamObserver#onNext(Object)} has completed. If disabled - * an application must make explicit calls to {@link #request} to receive messages. + * Disables automatic flow control, a mode where another message is implicitly {@link #request}ed + * after each call to the inbound's {@link StreamObserver#onNext(Object)} returns. + * + *

If disabled an application must make explicit calls to {@link #request} to receive messages. * *

On client-side this method may only be called during {@link * ClientResponseObserver#beforeStart}. On server-side it may only be called during the initial @@ -116,8 +134,7 @@ public abstract class CallStreamObserver implements StreamObserver { public abstract void disableAutoInboundFlowControl(); /** - * Requests the peer to produce {@code count} more messages to be delivered to the 'inbound' - * {@link StreamObserver}. + * Requests that {@code count} more messages be delivered to the 'inbound' {@link StreamObserver}. * *

This method is safe to call from multiple threads without external synchronization. * diff --git a/stub/src/main/java/io/grpc/stub/ClientCallStreamObserver.java b/stub/src/main/java/io/grpc/stub/ClientCallStreamObserver.java index 8f420fa77e4..4a173520359 100644 --- a/stub/src/main/java/io/grpc/stub/ClientCallStreamObserver.java +++ b/stub/src/main/java/io/grpc/stub/ClientCallStreamObserver.java @@ -92,8 +92,7 @@ public void disableAutoRequestWithInitial(int request) { public abstract void setOnReadyHandler(Runnable onReadyHandler); /** - * Requests the peer to produce {@code count} more messages to be delivered to the 'inbound' - * {@link StreamObserver}. + * Requests that {@code count} more messages be delivered to the 'inbound' {@link StreamObserver}. * *

This method is safe to call from multiple threads without external synchronization. * diff --git a/stub/src/main/java/io/grpc/stub/ServerCallStreamObserver.java b/stub/src/main/java/io/grpc/stub/ServerCallStreamObserver.java index 6ffea3500cc..cc5fbaf3fb0 100644 --- a/stub/src/main/java/io/grpc/stub/ServerCallStreamObserver.java +++ b/stub/src/main/java/io/grpc/stub/ServerCallStreamObserver.java @@ -147,8 +147,7 @@ public void disableAutoRequest() { public abstract void setOnReadyHandler(Runnable onReadyHandler); /** - * Requests the peer to produce {@code count} more messages to be delivered to the 'inbound' - * {@link StreamObserver}. + * Requests that {@code count} more messages be delivered to the 'inbound' {@link StreamObserver}. * *

This method is safe to call from multiple threads without external synchronization. *