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

> Master workerd streams with backpressure using its dual-stream architecture. Learn how to manage queue sizes and ensure flow control for efficient data handling.

- Repository: [Cloudflare/workerd](https://github.com/cloudflare/workerd)
- Tags: how-to-guide
- Published: 2026-03-18

---

**Workerd implements backpressure through a dual-stream architecture where standard JS-backed controllers in [`src/workerd/api/streams/standard.h`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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:

```cpp
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:

```javascript
// 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`](https://github.com/cloudflare/workerd/blob/main/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`:

```cpp
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:

```javascript
// 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`](https://github.com/cloudflare/workerd/blob/main/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

```javascript
// 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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/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`](https://github.com/cloudflare/workerd/blob/main/src/workerd/api/streams/readable.h) and [`writable.h`](https://github.com/cloudflare/workerd/blob/main/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.