Backpressure is the mechanism that lets a slow consumer slow its producer down instead of being buried. In a one-way stream there is one chain of producer, buffers and consumer to get right. A bidirectional stream has two chains running in opposite directions over the same connection, often through the same handler code, and every relay in the path doubles the number of places where a chain can break.

The basics, transport windows, demand signals and the menu of overflow policies, are covered in backpressure architecture in bidirectional streams. The classic send-everything-then-read deadlock is dissected in gRPC bidirectional streaming. This article picks up from there and focuses on propagation: how to carry a slow-down signal through a gateway that bridges two streams, how to code it with manual flow control in gRPC and with an explicit credit protocol over WebSocket, and how to prove with a test that it actually works.

Advertisement

Why propagation is the real problem

Every buffer in a streaming system answers one question: what happens when it fills? If the answer is "the writer waits", the slow-down signal moves one hop backward. If the answer is "it grows", the signal stops there and the buffer becomes a memory leak with a delay. Backpressure works only when every hop from the final consumer back to the original producer answers "the writer waits" or applies a deliberate drop policy.

Bidirectional streams make this harder in three ways. First, flow control at the transport is per direction: each receiver advertises its own windows, so the two directions never share one, and any coupling between them is created by application code. Second, that code is usually one handler serving both directions, and it is tempting to write it so that reading depends on writing. Third, real systems put relays in the middle: API gateways, WebSocket-to-gRPC bridges, service meshes and fan-out hubs. A relay terminates one stream and opens another, so the transport windows on its two sides are independent; only the relay's own code can connect them.

The four flows of a relay

Consider a speech product. A browser sends audio frames over a WebSocket to a gateway; the gateway opens a gRPC bidirectional stream to a recognition backend; transcripts flow back the same way. That is four flows: client to gateway (A), gateway to backend (B), backend to gateway (C) and gateway to client (D).

A relay bridging two bidi streams has four flows; each must pass demand backward, not buffer it awayBrowser clientWebSocketGateway / relayup queue (cap 32)client to backenddown queue (cap 32)backend to clientBackend servicegRPC bidi streamA: audio framesB: request(n)-pacedC: transcriptsD: credit-pacedUpstream demand: backend request(n) -> gateway reads WebSocket only when up queue has roomDownstream demand: client CREDIT frames -> gateway writes only with credit -> pauses gRPC readsBroken chainan unbounded queue, a fire-and-forget write, or a read loop that never pauses hides the stall until memory runs outEach direction has its own chain. The two chains must never wait on each other.
A gateway bridging a WebSocket to a gRPC bidi stream. Demand for the upstream chain comes from the backend; demand for the downstream chain comes from the browser. Each chain passes through one bounded queue in the gateway.

There are two demand chains, not one. Upstream, the backend's ability to consume audio must throttle how fast the gateway reads the WebSocket, which in turn must fill the client's TCP window and make the browser's send buffer grow, which the client must watch. Downstream, the browser's ability to render transcripts must throttle how fast the gateway writes to the WebSocket, which must throttle how fast it reads from gRPC, which must stop it granting HTTP/2 window to the backend, which must make the backend's sends block or report not-ready.

The rule that keeps the two chains from deadlocking each other is the same one the gRPC article gives for a single stream, applied per hop: the upstream chain and the downstream chain each run in their own task with their own bounded queue, and neither ever waits on the other. A gateway that reads one audio frame, forwards it, then waits for a transcript before reading the next frame has coupled the chains and will stall as soon as the backend batches its output.

Advertisement

Manual flow control in grpc-java

By default grpc-java requests inbound messages automatically, so a server handler receives messages as fast as the network delivers them regardless of whether it can process them. The library's manual flow control turns that off and gives you two levers: request(n) to ask for exactly n more inbound messages, and isReady() with setOnReadyHandler to learn whether outbound writes would currently be buffered beyond the transport's comfortable limit. The pattern below follows the manual flow control example in the grpc-java repository.

@Override
public StreamObserver<AudioChunk> recognize(StreamObserver<Transcript> responseObserver) {
  ServerCallStreamObserver<Transcript> out =
      (ServerCallStreamObserver<Transcript>) responseObserver;
  out.disableAutoRequest();                 // we decide when to read

  AtomicBoolean wasReady = new AtomicBoolean(false);
  out.setOnReadyHandler(() -> {
    // Called when the outbound side has room again. Resume reading only now.
    if (out.isReady() && wasReady.compareAndSet(false, true)) {
      out.request(1);
    }
  });

  return new StreamObserver<AudioChunk>() {
    @Override public void onNext(AudioChunk chunk) {
      for (Transcript t : recognizer.feed(chunk)) {
        out.onNext(t);                      // may buffer if the client is slow
      }
      if (out.isReady()) {
        out.request(1);                     // room downstream: take one more chunk
      } else {
        wasReady.set(false);                // pause; onReady will resume us
      }
    }
    @Override public void onError(Throwable t) { recognizer.close(); }
    @Override public void onCompleted() {
      for (Transcript t : recognizer.flush()) out.onNext(t);
      out.onCompleted();
    }
  };
}

This handler couples its input to its output deliberately, which is correct here because each audio chunk produces transcripts: if the client stops reading transcripts, the server stops reading audio, its HTTP/2 receive window stops being replenished, and the client's sends eventually report not-ready. That is propagation, not deadlock, as long as the client keeps its two directions in separate tasks. Coupling is safe on one side of a stream; coupling on both sides is the deadlock.

Two details are easy to miss. onNext on the response observer never blocks; isReady() is advisory, and writing while not ready simply grows an internal buffer, so a handler that ignores it has no backpressure at all. And the initial demand matters: the first request(1) here comes from the ready handler, which the library calls once the call starts. On the client side, ClientCallStreamObserver offers disableAutoRequestWithInitial(n) for the same purpose.

A credit protocol for WebSocket

The browser WebSocket API has no receive-side demand signal: messages arrive through onmessage as fast as the network delivers them, and the only send-side signal is bufferedAmount, which you have to poll. When the transport gives you nothing, build demand into the application protocol. The simplest scheme is credits: the receiver grants the sender a number of messages, the sender decrements on each send and stops at zero, and the receiver grants more as it finishes processing.

import asyncio, json

class CreditedSender:
    # Gateway side of flow D: send transcripts only while the browser has granted credit.
    def __init__(self, ws, queue_cap=32):
        self.ws, self.credit = ws, 0
        self.has_credit = asyncio.Event()
        self.queue = asyncio.Queue(maxsize=queue_cap)   # bounded: a full queue pauses the gRPC reader

    def on_control(self, msg):                          # called by the WebSocket reader task
        if msg.get("type") == "credit":
            self.credit += int(msg["n"])
            self.has_credit.set()

    async def pump(self):                               # one task per direction
        while True:
            item = await self.queue.get()
            while self.credit == 0:
                self.has_credit.clear()
                await asyncio.wait_for(self.has_credit.wait(), timeout=30)  # stalled client -> close
            self.credit -= 1
            await self.ws.send(json.dumps({"type": "transcript", "data": item}))

async def grpc_reader(call, sender):
    async for transcript in call:                       # stops pulling when the queue is full,
        await sender.queue.put(transcript)              # so gRPC stops granting window upstream

On the browser side the receiver starts by granting, say, 16 credits, and sends another grant of 8 each time it has rendered 8 messages. The grant size is a latency and memory trade-off: the number of credits outstanding must cover one round trip of messages at the target rate, or the sender idles waiting for the next grant. If the client renders 50 messages per second over an 80 ms round trip, at least 4 credits must be in flight; 16 gives headroom for jitter.

Notice the timeout. A client that stops granting is either slow or gone, and an open connection with zero credit holds a gateway queue indefinitely. Close it after a bound and let the client reconnect and resume, rather than holding resources for a peer that will never read.

Getting the queue sizes right

Every bounded queue in a relay should be sized from one number: how many messages need to be in flight to keep the next hop busy for one round trip. Larger queues do not increase throughput once that number is covered; they only add latency and memory, and they delay the moment the producer learns it must slow down. For the speech gateway, audio arrives at 50 frames per second; with a 40 ms round trip to the backend, two frames are enough to keep the backend fed, so a queue of 32 is generous and bounds the latency it can add to 640 ms.

Memory is the other bound. Multiply the per-connection queue capacity by the maximum message size and by the number of concurrent connections the gateway admits. If that product does not fit comfortably in the process, shrink the queues or admit fewer connections; load shedding for bidi streams covers refusing work at the front door when every chain is already saturated.

Proving it: the stall test

Backpressure bugs do not show up in functional tests, because test consumers are fast. Write a test that makes one consumer slow and asserts that the slowness reaches the producer within a bounded amount of buffered data. For the gateway, that means three runs.

  1. Downstream stall. A fake browser grants 4 credits and then stops. Assert that the gateway's down queue reaches its cap, that the backend's sends report not-ready within seconds, and that gateway memory stays flat while the backend keeps trying to produce.
  2. Upstream stall. A fake backend stops calling for more input. Assert that the gateway stops reading the WebSocket, and that the fake client's send buffer grows instead of the gateway's heap.
  3. Cross-direction independence. Stall downstream only, and assert that upstream audio still flows at full rate until the backend itself decides to pause. If upstream stops immediately, the gateway has coupled the two chains.

Measure buffered bytes, not message counts, and run each case for long enough that an unbounded buffer would be obvious, a minute is usually plenty. Keep the tests in CI: a single refactor that swaps a bounded queue for an unbounded one silently removes backpressure, and nothing else will catch it.

Failure modes

  • Hidden unbounded buffers. Library send queues, event-loop write buffers and logging pipelines can absorb a stall; check every hop, including ones you did not write.
  • Ignoring readiness. Writing to a gRPC observer without checking isReady() moves the backlog into the library's buffer.
  • Coupled directions in a relay. Waiting for a response before reading the next request stalls both chains when either side batches.
  • Connection-window starvation. On HTTP/2, one stream whose receiver stops reading can consume the shared connection window and starve other streams on the same connection; isolate heavy streams on separate connections if this bites.
  • Credit leaks. A lost or never-sent credit grant leaves the sender at zero forever; reconnect logic must reset credit on both sides.
  • Silent stalls. A stream stuck at zero credit or zero window produces no errors; without per-stream stall metrics and timeouts it looks healthy.

Operations and trade-offs

Export, per connection and per direction: queue depth, time spent with zero credit or not-ready, and bytes buffered. Alert on the fraction of streams stalled longer than a threshold, not on any single stall. Tight queues give fast propagation and low latency but make throughput sensitive to jitter; loose queues smooth jitter but delay the signal and cost memory. Blocking propagation is right for data that must not be lost, such as audio for transcription; for state feeds, conflating to the latest value at the relay is often better than slowing the producer. For transport details see the WebSocket guide.

What to do next

  1. Draw every hop of one bidi path, both directions, and label each buffer with its bound and what happens when it fills.
  2. Replace every unbounded queue in your relays with a bounded one sized from rate times round-trip time.
  3. Run the upstream and downstream chains in separate tasks and remove any read-then-wait-for-response logic.
  4. Switch gRPC handlers to manual flow control with request(n) and honour isReady() on every write.
  5. Add a credit protocol, with a stall timeout, to any WebSocket path that carries more than the client can render.
  6. Write the three stall tests and run them in CI.
  7. Export per-direction stall time and queue depth, and alert on the fraction of stalled streams.
Key takeaway: A bidirectional stream carries two independent demand chains, and every relay in the path must connect each chain across its two sides with a bounded queue. Use manual flow control in gRPC so reads follow downstream readiness. Add an explicit credit protocol where the transport has none, as with browser WebSockets. Keep the two directions in separate tasks so they never wait on each other. Then prove propagation with stall tests, because no functional test will ever show you a missing link.