How to Implement Streams with Backpressure in Workerd: A Complete Guide

Workerd implements backpressure through a dual-stream architecture where standard JS-backed controllers in src/workerd/api/streams/standard.h track queue size against configurable high-water marks, exposing flow control via controller.desiredSize and writer.ready promises while halting the pull algorithm when buffers exceed capacity.

Workert is Cloudflare’s open-source JavaScript runtime that powers Workers, implementing the WHATWG Streams specification with high-performance C++ primitives. Understanding how to implement streams with backpressure in workerd requires navigating its two distinct stream implementations and the specific queue-based mechanisms that throttle data producers when consumers slow down. The runtime manages flow control by comparing buffered data against high-water marks in the ReadableImpl and WritableImpl controllers, ensuring memory remains bounded even with slow consumers.

Understanding Workerd's Dual Stream Architecture

Workerd implements the Streams specification twice to optimize for different execution contexts. The repository maintains internal streams—thin wrappers around Cloudflare’s kj::Async* I/O primitives defined in headers like src/workerd/api/streams/readable.h—and standard JS-backed streams that provide full spec compliance within the V8 isolate.

Internal streams handle byte data exclusively and run outside the V8 isolate lock, delegating backpressure to the underlying kj queue. Standard streams operate inside the isolate lock and rely on the StateMachine class from src/workerd/util/state-machine.h alongside queue abstractions (ValueQueue for value streams, ByteQueue for byte streams). Both implementations track how many units are buffered; when this count exceeds the high-water mark (HWM), the stream signals backpressure to pause production.

Implementing Backpressure in Readable Streams

In the standard implementation located in src/workerd/api/streams/standard.h, the ReadableStreamDefaultController owns a ReadableImpl object (lines 31-45) that manages the pull lifecycle and queue state. The controller computes the desired size—the number of chunks that can be enqueued before reaching the HWM—to determine whether to request additional data from the underlying source.

When desiredSize drops to zero or below, the shouldCallPull() method (lines 79-81) returns false, preventing pullIfNeeded() (lines 71-74) from scheduling further pull callbacks until the consumer drains the queue:

kj::Maybe<int> ReadableImpl::getDesiredSize() {
  // HWM – current buffered units
  int size = highWaterMark - queue.size();
  return size < 0 ? kj::none : kj::some(size);
}

void ReadableImpl::pullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) {
  if (shouldCallPull()) {
    // invoke the underlying source’s pull algorithm
    algorithms.pull->call(js, ...);
    flags.pulling = true;
  }
}

This mechanism effectively propagates backpressure upstream: once the queue reaches the HWM, no further pull callbacks are invoked until desiredSize becomes positive again.

Creating a Back-Pressured ReadableStream

Configure the highWaterMark in the constructor options and monitor desiredSize through the reader to observe backpressure in action:

// src/example/backpressured-readable.js
const highWaterMark = 2;               // HWM = 2 chunks
let counter = 0;

const readable = new ReadableStream(
  {
    async pull(controller) {
      // Called only while queued chunks < highWaterMark
      counter++;
      controller.enqueue(`msg-${counter}`);
    }
  },
  { highWaterMark }                    // <- passed to the C++ controller
);

(async () => {
  const reader = readable.getReader();

  // Desired size starts at 2, drops to 1 after first enqueue, 0 after second.
  console.log('desiredSize →', reader.desiredSize); // 2

  // First read triggers pull (queue empty) -> pull enqueues msg‑1
  console.log(await reader.read()); // {value: "msg-1", done: false}
  console.log('desiredSize →', reader.desiredSize); // 1

  // Second read → pull enqueues msg‑2 (queue now full, desiredSize 0)
  console.log(await reader.read()); // {value: "msg-2", done: false}
  console.log('desiredSize →', reader.desiredSize); // 0

  // Third read: queue empty, pull is *not* called until the consumer reads again,
  // because desiredSize is 0 (backpressure). Pull will fire when we request.
  console.log(await reader.read()); // {value: "msg-3", done: false}
})();

Implementing Backpressure in Writable Streams

The writable side tracks backpressure through the WritableImpl class (lines 69-84 in standard.h). This implementation maintains amountBuffered—the sum of sizes of all queued write requests—and compares it against the highWaterMark.

The getDesiredSize() method (lines 19-21) returns highWaterMark - amountBuffered. When this value drops to zero or below, the updateBackpressure() method (lines 36-40) sets the internal backpressure flag and creates a new promise for writer.ready:

ssize_t WritableImpl::getDesiredSize() {
  return static_cast<ssize_t>(highWaterMark) - static_cast<ssize_t>(amountBuffered);
}

void WritableImpl::updateBackpressure(jsg::Lock& js) {
  bool nowBackpressured = getDesiredSize() <= 0;
  if (nowBackpressured != flags.backpressure) {
    flags.backpressure = nowBackpressured;
    if (nowBackpressured) {
      // a new promise for writer.ready
      ready = jsg::PromiseResolverPair<void>();
    } else {
      // resolve the previous ready promise
      ready.resolve();
    }
  }
}

JavaScript code observes backpressure through two mechanisms: writer.desiredSize (a numeric hint, or null when over-full) and writer.ready (a promise that resolves when pressure clears). When desiredSize is less than or equal to zero, the producer should await writer.ready before enqueueing additional data.

Implementing a WritableStream that Respects Backpressure

Set a low highWaterMark to constrain memory usage and await writer.ready before submitting additional chunks:

// src/example/backpressured-writable.js
const highWaterMark = 1;                // only one pending write allowed

const writable = new WritableStream({
  async write(chunk, controller) {
    // Simulate a slow sink
    await new Promise(r => setTimeout(r, 50));
    console.log('written:', chunk);
  }
}, { highWaterMark });

(async () => {
  const writer = writable.getWriter();

  console.log('desiredSize →', writer.desiredSize); // 1

  // First write succeeds immediately, desiredSize becomes 0
  writer.write('A');
  console.log('desiredSize after A →', writer.desiredSize); // 0

  // Second write is queued; backpressure kicks in (ready becomes a new pending promise)
  const backpressPromise = writer.ready;
  writer.write('B');
  console.log('desiredSize after B →', writer.desiredSize); // -1 (over‑full)

  // Wait for backpressure to clear (ready resolves)
  await backpressPromise;
  console.log('backpressure cleared, desiredSize →', writer.desiredSize); // 0 or 1

  await writer.close();
})();

Piping with Backpressure Across Stream Types

When piping streams in workerd, the runtime selects from four specialized loops depending on the source and sink types, as documented in docs/streams.md. The most efficient path—kj-to-kj—bypasses V8 entirely, letting kernel-level kj::AsyncInputStream write directly into kj::AsyncOutputStream. The JS-to-JS pipe loop honors backpressure through the mechanisms described above, propagating signals from the writable sink back to the readable source.

The pipeTo() method automatically pauses pulls when the sink's queue fills and resumes when writer.ready resolves, preventing unbounded memory growth during pipe operations.

Pipe Between Streams with Backpressure

// src/example/pipe-backpressure.js
const source = new ReadableStream({
  async pull(controller) {
    // Produce data slowly to let backpressure manifest
    await new Promise(r => setTimeout(r, 30));
    controller.enqueue(Math.random());
    if (Math.random() < 0.1) controller.close();
  }
}, { highWaterMark: 4 });

const sink = new WritableStream({
  async write(chunk) {
    // Slow consumer – backpressure will propagate upstream
    await new Promise(r => setTimeout(r, 100));
    console.log('sink received', chunk);
  }
}, { highWaterMark: 2 });

await source.pipeTo(sink);   // Internally uses the JS‑to‑JS pipe loop.

Summary

  • Workerd implements two stream types: internal kj-based streams for I/O outside the isolate lock, and standard JS-backed streams that run within V8 and expose spec-compliant backpressure signals through src/workerd/api/streams/standard.h.
  • High-water marks control flow: Both ReadableImpl and WritableImpl compare buffered data against the HWM to compute desiredSize, with default values of 1 for both stream types.
  • Backpressure signals differ by direction: Readable streams stop calling pull when desiredSize reaches zero; writable streams set the backpressure flag and reset writer.ready when the queue fills.
  • State machines manage transitions: The StateMachine utility in src/workerd/util/state-machine.h coordinates the underlying state transitions for both readable and writable controllers.
  • Cross-type piping works automatically: The pipeTo() implementation selects between four internal loops, with JS-to-JS pipes correctly propagating backpressure through the ready promise and pull scheduling.

Frequently Asked Questions

How does workerd differ from Node.js streams for backpressure?

Workerd follows the WHATWG Streams specification strictly, using desiredSize and writer.ready promises for flow control, whereas Node.js uses the legacy write() return value and 'drain' events. The implementation in src/workerd/api/streams/standard.h runs inside the V8 isolate lock and relies on C++ ReadableImpl and WritableImpl templates to manage queue state, unlike Node.js's C++ stream wrappers which operate outside the JavaScript engine's control.

What is the default high-water mark in workerd streams?

The default high-water mark is 1 for both readable and writable streams, as defined in the ReadableImpl and WritableImpl constructors in src/workerd/api/streams/standard.h. You can override this in the stream constructor options to buffer more data and reduce context switching, though higher values increase memory usage within the isolate.

Can I disable backpressure in workerd streams?

You cannot fully disable backpressure in standard streams because the queue management is intrinsic to the StateMachine implementation in src/workerd/util/state-machine.h. However, you can effectively disable flow control by setting an extremely high highWaterMark value, though this risks unbounded memory growth if the producer outpaces the consumer.

How do internal streams handle backpressure differently?

Internal streams (defined in headers like src/workerd/api/streams/readable.h and writable.h) are thin wrappers around kj::Async* primitives that run outside the V8 isolate lock. They rely on the underlying kj queue for backpressure rather than the JavaScript-visible desiredSize mechanism, making them more efficient for byte-only I/O but less flexible for JavaScript-based flow control logic.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →