Skip to content
File

Blob: src/workerd/api/tests/streams-backpressure-test.js

javascript367 lines
1// Copyright (c) 2026 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5// Tests for backpressure behavior in JS-backed streams.
6// These tests focus on desiredSize tracking, ready promise behavior,
7// and backpressure propagation through pipe chains.
8//
9// Test inspirations:
10// - Deno: tests/unit/streams_test.ts (parameterized count/delay tests)
11// - Deno: tests/node_compat/test-stream-readable-hwm-0.js (hwm=0 edge case)
12// - Bun: test/js/node/stream/node-stream.test.js (backpressure tests)
13// - Bun: test/js/node/http/node-http-backpressure.test.ts (HTTP-level backpressure)
14// - Bun: test/js/web/fetch/body-stream.test.ts (backpressure with various input lengths)
15 
16import { strictEqual, ok } from 'node:assert';
17 
18// Test ReadableStream with highWaterMark = 0
19// Pull should only be called when there's an active read request
20// Inspired by: Deno tests/node_compat/test-stream-readable-hwm-0.js (hwm=0 edge case)
21export const backpressureReadableHwmZero = {
22 async test() {
23 let pullCount = 0;
24 let controller;
25 
26 const rs = new ReadableStream(
27 {
28 start(c) {
29 controller = c;
30 },
31 pull(c) {
32 pullCount++;
33 c.enqueue(pullCount);
34 },
35 },
36 { highWaterMark: 0 }
37 );
38 
39 await scheduler.wait(10);
40 strictEqual(pullCount, 0, 'pull should not be called without read');
41 strictEqual(controller.desiredSize, 0, 'desiredSize should be 0');
42 
43 const reader = rs.getReader();
44 
45 const read1 = reader.read();
46 await scheduler.wait(1);
47 strictEqual(pullCount, 1, 'pull should be called once after read');
48 
49 const result1 = await read1;
50 strictEqual(result1.value, 1);
51 strictEqual(result1.done, false);
52 
53 strictEqual(controller.desiredSize, 0, 'desiredSize should remain 0');
54 
55 const read2 = reader.read();
56 await scheduler.wait(1);
57 strictEqual(pullCount, 2, 'pull should be called again');
58 
59 const result2 = await read2;
60 strictEqual(result2.value, 2);
61 
62 reader.releaseLock();
63 },
64};
65 
66// Test ReadableStream with highWaterMark = 1
67// Verify desiredSize transitions correctly
68// Inspired by: Deno tests/unit/streams_test.ts (parameterized hwm tests)
69export const backpressureReadableHwmOne = {
70 async test() {
71 let pullCount = 0;
72 
73 const rs = new ReadableStream(
74 {
75 start(c) {
76 strictEqual(c.desiredSize, 1, 'initial desiredSize should be 1');
77 },
78 pull(c) {
79 pullCount++;
80 c.enqueue(pullCount);
81 },
82 },
83 { highWaterMark: 1 }
84 );
85 
86 const reader = rs.getReader();
87 
88 const result1 = await reader.read();
89 strictEqual(result1.value, 1);
90 strictEqual(result1.done, false);
91 
92 const result2 = await reader.read();
93 strictEqual(result2.value, 2);
94 strictEqual(result2.done, false);
95 
96 reader.releaseLock();
97 },
98};
99 
100// Test ReadableStream with highWaterMark = 64
101// Multiple chunks should be buffered
102// Inspired by: Bun test/js/web/streams/streams.test.js (large buffer tests)
103export const backpressureReadableHwmLarge = {
104 async test() {
105 let pullCount = 0;
106 const MAX_PULLS = 100;
107 
108 const rs = new ReadableStream(
109 {
110 pull(c) {
111 pullCount++;
112 if (pullCount <= MAX_PULLS) {
113 c.enqueue(pullCount);
114 } else {
115 c.close();
116 }
117 },
118 },
119 { highWaterMark: 64 }
120 );
121 
122 const reader = rs.getReader();
123 
124 const values = [];
125 while (true) {
126 const { value, done } = await reader.read();
127 if (done) break;
128 values.push(value);
129 }
130 
131 strictEqual(values.length, MAX_PULLS);
132 strictEqual(values[0], 1);
133 strictEqual(values[99], 100);
134 ok(pullCount > MAX_PULLS, 'pull called enough times');
135 },
136};
137 
138// Test WritableStream desiredSize tracking and ready promise
139// Inspired by: Bun test/js/web/fetch/body-stream.test.ts (backpressure with various input lengths)
140export const backpressureWritableDesiredSize = {
141 async test() {
142 const written = [];
143 
144 const ws = new WritableStream(
145 {
146 write(chunk) {
147 written.push(chunk);
148 // Simulate slow write
149 return new Promise((resolve) => setTimeout(resolve, 10));
150 },
151 },
152 { highWaterMark: 3 }
153 );
154 
155 const writer = ws.getWriter();
156 
157 strictEqual(writer.desiredSize, 3, 'initial desiredSize');
158 
159 writer.write(1);
160 strictEqual(writer.desiredSize, 2, 'desiredSize after first write');
161 
162 writer.write(2);
163 strictEqual(writer.desiredSize, 1, 'desiredSize after second write');
164 
165 writer.write(3);
166 strictEqual(writer.desiredSize, 0, 'desiredSize at capacity');
167 
168 const readyBefore = writer.ready;
169 ok(readyBefore instanceof Promise, 'ready should be a Promise');
170 
171 writer.write(4);
172 strictEqual(writer.desiredSize, -1, 'desiredSize over capacity');
173 
174 const readyAfter = writer.ready;
175 ok(readyAfter instanceof Promise, 'ready should still be a Promise');
176 
177 await scheduler.wait(50);
178 
179 ok(writer.desiredSize > -1, 'desiredSize recovered after writes complete');
180 
181 const readyResolved = writer.ready;
182 await readyResolved;
183 
184 await writer.close();
185 strictEqual(written.length, 4, 'all chunks written');
186 },
187};
188 
189// Test WritableStream with slow sink causes backpressure
190// Inspired by: Bun test/js/bun/spawn/spawn-stdin-readable-stream-edge-cases.test.ts (slow consumer)
191export const backpressureWritableSlowSink = {
192 async test() {
193 let writeCount = 0;
194 
195 const ws = new WritableStream(
196 {
197 async write() {
198 writeCount++;
199 // Very slow write
200 await scheduler.wait(50);
201 },
202 },
203 { highWaterMark: 2 }
204 );
205 
206 const writer = ws.getWriter();
207 
208 strictEqual(writer.desiredSize, 2);
209 
210 const w1 = writer.write('a');
211 strictEqual(writer.desiredSize, 1);
212 
213 const w2 = writer.write('b');
214 strictEqual(writer.desiredSize, 0);
215 
216 const initialReady = writer.ready;
217 
218 const w3 = writer.write('c');
219 strictEqual(writer.desiredSize, -1);
220 
221 // Backpressure can be signaled in two ways depending on implementation:
222 // 1. writer.ready returns a new pending promise (ready !== initialReady)
223 // 2. desiredSize drops to 0 or below, indicating the queue is full
224 // We accept either signal as valid backpressure indication.
225 ok(
226 writer.ready !== initialReady || writer.desiredSize <= 0,
227 'backpressure signal'
228 );
229 
230 const results = await Promise.allSettled([w1, w2, w3]);
231 strictEqual(results[0].status, 'fulfilled', 'first write should succeed');
232 strictEqual(results[1].status, 'fulfilled', 'second write should succeed');
233 strictEqual(results[2].status, 'fulfilled', 'third write should succeed');
234 strictEqual(writeCount, 3);
235 
236 await writer.close();
237 },
238};
239 
240// Test TransformStream with both readable and writable strategies
241// Inspired by: Bun test/js/web/streams/streams.test.js (TransformStream tests)
242export const backpressureTransformBothStrategies = {
243 async test() {
244 const ts = new TransformStream(
245 {
246 transform(chunk, controller) {
247 controller.enqueue(chunk);
248 },
249 },
250 { highWaterMark: 2 },
251 { highWaterMark: 4 }
252 );
253 
254 const writer = ts.writable.getWriter();
255 const reader = ts.readable.getReader();
256 
257 strictEqual(writer.desiredSize, 2, 'writable desiredSize');
258 
259 writer.write(1);
260 writer.write(2);
261 writer.write(3);
262 writer.write(4);
263 
264 const results = [];
265 for (let i = 0; i < 4; i++) {
266 const { value } = await reader.read();
267 results.push(value);
268 }
269 
270 strictEqual(results.length, 4);
271 strictEqual(results[0], 1);
272 strictEqual(results[3], 4);
273 
274 await writer.close();
275 const final = await reader.read();
276 ok(final.done);
277 },
278};
279 
280// Test backpressure propagation through pipe chain
281// Inspired by: Bun test/js/web/streams/streams.test.js (pipeTo/pipeThrough tests)
282export const backpressurePipeChain = {
283 async test() {
284 let sourcePullCount = 0;
285 const MAX_CHUNKS = 5;
286 
287 const source = new ReadableStream(
288 {
289 pull(c) {
290 sourcePullCount++;
291 if (sourcePullCount <= MAX_CHUNKS) {
292 c.enqueue(sourcePullCount);
293 } else {
294 c.close();
295 }
296 },
297 },
298 { highWaterMark: 2 }
299 );
300 
301 const transform = new TransformStream(
302 {
303 transform(chunk, controller) {
304 controller.enqueue(chunk * 2);
305 },
306 },
307 { highWaterMark: 1 },
308 { highWaterMark: 1 }
309 );
310 
311 const chunks = [];
312 const dest = new WritableStream(
313 {
314 async write(chunk) {
315 chunks.push(chunk);
316 // Slow consumer
317 await scheduler.wait(5);
318 },
319 },
320 { highWaterMark: 1 }
321 );
322 
323 await source.pipeThrough(transform).pipeTo(dest);
324 
325 strictEqual(chunks.length, MAX_CHUNKS, 'all chunks received');
326 strictEqual(chunks[0], 2, 'first chunk transformed');
327 strictEqual(chunks[4], 10, 'last chunk transformed');
328 },
329};
330 
331// Test byte stream highWaterMark is measured in bytes, not chunks
332// Inspired by: Bun test/js/node/test/parallel/test-whatwg-readablebytestream.js
333export const backpressureByteStreamHwm = {
334 async test() {
335 let controller;
336 let pullCount = 0;
337 
338 const rs = new ReadableStream(
339 {
340 type: 'bytes',
341 start(c) {
342 controller = c;
343 strictEqual(c.desiredSize, 10, 'initial desiredSize in bytes');
344 },
345 pull(c) {
346 pullCount++;
347 // Enqueue 3 bytes
348 c.enqueue(new Uint8Array([1, 2, 3]));
349 },
350 },
351 { highWaterMark: 10 }
352 );
353 
354 await scheduler.wait(20);
355 
356 ok(pullCount >= 3, 'pulled multiple times for byte count');
357 ok(controller.desiredSize <= 1, 'desiredSize accounts for bytes');
358 
359 const reader = rs.getReader();
360 
361 const { value } = await reader.read();
362 ok(value.byteLength >= 9, 'received buffered bytes');
363 
364 reader.releaseLock();
365 },
366};