Skip to content
File

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

javascript537 lines
1// Copyright (c) 2023 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 
5import * as assert from 'node:assert';
6import * as util from 'node:util';
7 
8export const partiallyReadStream = {
9 async test(ctrl, env, ctx) {
10 const enc = new TextEncoder();
11 const rs = new ReadableStream({
12 type: 'bytes',
13 start(controller) {
14 controller.enqueue(enc.encode('hello'));
15 controller.enqueue(enc.encode('world'));
16 controller.close();
17 },
18 });
19 const reader = rs.getReader({ mode: 'byob' });
20 await reader.read(new Uint8Array(5));
21 reader.releaseLock();
22 
23 // Should not throw!
24 await env.KV.put('key', rs);
25 },
26};
27 
28export const arrayBufferOfReadable = {
29 async test() {
30 const cs = new CompressionStream('gzip');
31 const cw = cs.writable.getWriter();
32 await cw.write(new TextEncoder().encode('0123456789'.repeat(1000)));
33 await cw.close();
34 const data = await new Response(cs.readable).arrayBuffer();
35 assert.equal(66, data.byteLength);
36 
37 const ds = new DecompressionStream('gzip');
38 const dw = ds.writable.getWriter();
39 await dw.write(data);
40 await dw.close();
41 
42 const read = await new Response(ds.readable).arrayBuffer();
43 assert.equal(10_000, read.byteLength);
44 },
45};
46 
47export const inspect = {
48 async test() {
49 const inspectOpts = { breakLength: Infinity };
50 
51 // Check with JavaScript regular ReadableStream
52 {
53 let pulls = 0;
54 const readableStream = new ReadableStream({
55 pull(controller) {
56 if (pulls === 0) controller.enqueue('hello');
57 if (pulls === 1) controller.close();
58 pulls++;
59 },
60 });
61 assert.strictEqual(
62 util.inspect(readableStream, inspectOpts),
63 "ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: false, [length]: undefined }"
64 );
65 
66 const reader = readableStream.getReader();
67 assert.strictEqual(
68 util.inspect(readableStream, inspectOpts),
69 "ReadableStream { locked: true, [state]: 'readable', [supportsBYOB]: false, [length]: undefined }"
70 );
71 
72 await reader.read();
73 assert.strictEqual(
74 util.inspect(readableStream, inspectOpts),
75 "ReadableStream { locked: true, [state]: 'readable', [supportsBYOB]: false, [length]: undefined }"
76 );
77 
78 await reader.read();
79 assert.strictEqual(
80 util.inspect(readableStream, inspectOpts),
81 "ReadableStream { locked: true, [state]: 'closed', [supportsBYOB]: false, [length]: undefined }"
82 );
83 }
84 
85 // Check with errored JavaScript regular ReadableStream
86 {
87 const readableStream = new ReadableStream({
88 start(controller) {
89 controller.error(new Error('Oops!'));
90 },
91 });
92 assert.strictEqual(
93 util.inspect(readableStream, inspectOpts),
94 "ReadableStream { locked: false, [state]: 'errored', [supportsBYOB]: false, [length]: undefined }"
95 );
96 }
97 
98 // Check with JavaScript bytes ReadableStream
99 {
100 const readableStream = new ReadableStream({
101 type: 'bytes',
102 pull(controller) {
103 controller.enqueue(new Uint8Array([1]));
104 },
105 });
106 assert.strictEqual(
107 util.inspect(readableStream, inspectOpts),
108 "ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: undefined }"
109 );
110 }
111 
112 // Check with JavaScript WritableStream
113 {
114 const writableStream = new WritableStream({
115 write(chunk, controller) {},
116 });
117 assert.strictEqual(
118 util.inspect(writableStream, inspectOpts),
119 "WritableStream { locked: false, [state]: 'writable', [expectsBytes]: false }"
120 );
121 
122 const writer = writableStream.getWriter();
123 assert.strictEqual(
124 util.inspect(writableStream, inspectOpts),
125 "WritableStream { locked: true, [state]: 'writable', [expectsBytes]: false }"
126 );
127 
128 await writer.write('chunk');
129 assert.strictEqual(
130 util.inspect(writableStream, inspectOpts),
131 "WritableStream { locked: true, [state]: 'writable', [expectsBytes]: false }"
132 );
133 
134 await writer.close();
135 assert.strictEqual(
136 util.inspect(writableStream, inspectOpts),
137 "WritableStream { locked: true, [state]: 'closed', [expectsBytes]: false }"
138 );
139 }
140 
141 // Check with errored JavaScript WritableStream
142 {
143 const writableStream = new WritableStream({
144 write(chunk, controller) {
145 controller.error(new Error('Oops!'));
146 },
147 });
148 assert.strictEqual(
149 util.inspect(writableStream, inspectOpts),
150 "WritableStream { locked: false, [state]: 'writable', [expectsBytes]: false }"
151 );
152 
153 const writer = writableStream.getWriter();
154 const promise = writer.write('chunk');
155 assert.strictEqual(
156 util.inspect(writableStream, inspectOpts),
157 "WritableStream { locked: true, [state]: 'erroring', [expectsBytes]: false }"
158 );
159 
160 await promise;
161 assert.strictEqual(
162 util.inspect(writableStream, inspectOpts),
163 "WritableStream { locked: true, [state]: 'errored', [expectsBytes]: false }"
164 );
165 }
166 
167 // Check with internal known-length TransformStream
168 {
169 const inspectOpts = { breakLength: 100 };
170 const transformStream = new FixedLengthStream(5);
171 assert.strictEqual(
172 util.inspect(transformStream, inspectOpts),
173 `FixedLengthStream {
174 readable: ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: 5n },
175 writable: WritableStream { locked: false, [state]: 'writable', [expectsBytes]: true }
176}`
177 );
178 
179 const { writable, readable } = transformStream;
180 const writer = writable.getWriter();
181 assert.strictEqual(
182 util.inspect(transformStream, inspectOpts),
183 `FixedLengthStream {
184 readable: ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: 5n },
185 writable: WritableStream { locked: true, [state]: 'writable', [expectsBytes]: true }
186}`
187 );
188 
189 void writer.write(new Uint8Array([1, 2, 3]));
190 void writer.write(new Uint8Array([4, 5]));
191 assert.strictEqual(
192 util.inspect(transformStream, inspectOpts),
193 `FixedLengthStream {
194 readable: ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: 5n },
195 writable: WritableStream { locked: true, [state]: 'writable', [expectsBytes]: true }
196}`
197 );
198 
199 void writer.close();
200 assert.strictEqual(
201 util.inspect(transformStream, inspectOpts),
202 `FixedLengthStream {
203 readable: ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: 5n },
204 writable: WritableStream { locked: true, [state]: 'closed', [expectsBytes]: true }
205}`
206 );
207 
208 const reader = readable.getReader();
209 assert.strictEqual(
210 util.inspect(transformStream, inspectOpts),
211 `FixedLengthStream {
212 readable: ReadableStream { locked: true, [state]: 'readable', [supportsBYOB]: true, [length]: 5n },
213 writable: WritableStream { locked: true, [state]: 'closed', [expectsBytes]: true }
214}`
215 );
216 
217 await reader.read();
218 assert.strictEqual(
219 util.inspect(transformStream, inspectOpts),
220 `FixedLengthStream {
221 readable: ReadableStream { locked: true, [state]: 'readable', [supportsBYOB]: true, [length]: 2n },
222 writable: WritableStream { locked: true, [state]: 'closed', [expectsBytes]: true }
223}`
224 );
225 
226 await reader.read();
227 assert.strictEqual(
228 util.inspect(transformStream, inspectOpts),
229 `FixedLengthStream {
230 readable: ReadableStream { locked: true, [state]: 'readable', [supportsBYOB]: true, [length]: 0n },
231 writable: WritableStream { locked: true, [state]: 'closed', [expectsBytes]: true }
232}`
233 );
234 
235 await reader.read();
236 assert.strictEqual(
237 util.inspect(transformStream, inspectOpts),
238 `FixedLengthStream {
239 readable: ReadableStream { locked: true, [state]: 'closed', [supportsBYOB]: true, [length]: 0n },
240 writable: WritableStream { locked: true, [state]: 'closed', [expectsBytes]: true }
241}`
242 );
243 }
244 
245 // Check with errored internal TransformStream
246 {
247 const inspectOpts = { breakLength: 100 };
248 const transformStream = new IdentityTransformStream();
249 assert.strictEqual(
250 util.inspect(transformStream, inspectOpts),
251 `IdentityTransformStream {
252 readable: ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: undefined },
253 writable: WritableStream { locked: false, [state]: 'writable', [expectsBytes]: true }
254}`
255 );
256 
257 const { writable, readable } = transformStream;
258 const writer = writable.getWriter();
259 void writer.abort(new Error('Oops!'));
260 assert.strictEqual(
261 util.inspect(transformStream, inspectOpts),
262 `IdentityTransformStream {
263 readable: ReadableStream { locked: false, [state]: 'readable', [supportsBYOB]: true, [length]: undefined },
264 writable: WritableStream { locked: true, [state]: 'errored', [expectsBytes]: true }
265}`
266 );
267 
268 const reader = readable.getReader();
269 assert.strictEqual(
270 util.inspect(transformStream, inspectOpts),
271 `IdentityTransformStream {
272 readable: ReadableStream { locked: true, [state]: 'readable', [supportsBYOB]: true, [length]: undefined },
273 writable: WritableStream { locked: true, [state]: 'errored', [expectsBytes]: true }
274}`
275 );
276 
277 await reader.read().catch(() => {});
278 assert.strictEqual(
279 util.inspect(transformStream, inspectOpts),
280 `IdentityTransformStream {
281 readable: ReadableStream { locked: true, [state]: 'errored', [supportsBYOB]: true, [length]: undefined },
282 writable: WritableStream { locked: true, [state]: 'errored', [expectsBytes]: true }
283}`
284 );
285 }
286 },
287};
288 
289// Test for re-entrancy bug: when pushing to multiple consumers (via tee),
290// the transform function can directly cancel another consumer synchronously.
291// This should not crash - the cancelled consumer should gracefully ignore the push.
292// Before the fix, this would crash with:
293// "expected state.template tryGet<Ready>() != nullptr; The consumer is either closed or errored."
294//
295// This test simulates the production scenario where:
296// 1. A TransformStream's readable is tee'd
297// 2. The transform function synchronously cancels one of the tee branches
298// 3. When enqueue is called, the push loop tries to push to the cancelled consumer
299export const transformTeeReentrancySynchronousCancel = {
300 async test() {
301 let reader2;
302 let cancelledBranch2 = false;
303 
304 // Create a TransformStream whose transform function cancels branch2
305 const ts = new TransformStream({
306 transform(chunk, controller) {
307 // First time through, cancel branch2 BEFORE enqueueing
308 // This simulates the production scenario where user code in the
309 // transform function affects another consumer
310 if (!cancelledBranch2 && reader2) {
311 reader2.cancel('cancelled synchronously in transform');
312 cancelledBranch2 = true;
313 }
314 controller.enqueue(chunk);
315 },
316 });
317 
318 const writer = ts.writable.getWriter();
319 const [branch1, branch2] = ts.readable.tee();
320 const reader1 = branch1.getReader();
321 reader2 = branch2.getReader();
322 
323 // Start pending reads on both branches
324 const read1Promise = reader1.read();
325 const read2Promise = reader2.read();
326 
327 // Write to the transform - this triggers the transform function which:
328 // 1. Cancels branch2 (closing/erroring its consumer)
329 // 2. Calls controller.enqueue() which pushes to all consumers
330 // Before the fix, step 2 would crash when trying to push to cancelled branch2
331 await writer.write('test data');
332 
333 // Verify branch1 got the data
334 const result1 = await read1Promise;
335 assert.strictEqual(result1.value, 'test data');
336 
337 // branch2 was cancelled, so its read should complete (done or with data before cancel)
338 const result2 = await read2Promise;
339 assert.ok(result2 !== undefined);
340 
341 await writer.close();
342 await reader1.cancel();
343 },
344};
345 
346// Test with TransformStream to match the production stack trace more closely.
347// The production bug occurred during: TransformStream โ†’ enqueue โ†’ QueueImpl::push โ†’ consumer iteration
348export const transformStreamTeeReentrancy = {
349 async test() {
350 const { readable, writable } = new TransformStream();
351 const writer = writable.getWriter();
352 
353 const [branch1, branch2] = readable.tee();
354 const reader1 = branch1.getReader();
355 const reader2 = branch2.getReader();
356 
357 // Start pending reads on both branches
358 const read1Promise = reader1.read();
359 const read2Promise = reader2.read();
360 
361 // When read1 resolves, cancel branch2
362 // This simulates the production scenario where a .then() handler
363 // attached to the read promise cancels another branch
364 read1Promise.then(() => {
365 reader2.cancel('cancelled during transform');
366 });
367 
368 // Write through the transform - this triggers enqueue on the readable side.
369 // Before the fix, this would crash when the push loop tried to push to
370 // the now-cancelled branch2 consumer.
371 await writer.write('transform data');
372 
373 // Verify read1 succeeded
374 const result1 = await read1Promise;
375 assert.strictEqual(result1.value, 'transform data');
376 
377 // read2 may have received data or be done - either is fine
378 // The important thing is no crash occurred
379 const result2 = await read2Promise;
380 assert.ok(result2 !== undefined);
381 
382 await writer.close();
383 await reader1.cancel();
384 },
385};
386 
387// Test that multiple writes through a tee'd ReadableStream work correctly
388// even when one branch is cancelled mid-stream
389export const teeWithCancelMidStream = {
390 async test() {
391 let controller;
392 const stream = new ReadableStream({
393 start(c) {
394 controller = c;
395 },
396 });
397 
398 const [branch1, branch2] = stream.tee();
399 const reader1 = branch1.getReader();
400 const reader2 = branch2.getReader();
401 
402 // Start reads on both branches
403 let read1Promise = reader1.read();
404 let read2Promise = reader2.read();
405 
406 // Enqueue first chunk - both branches should get it
407 controller.enqueue('chunk1');
408 const r1a = await read1Promise;
409 const r2a = await read2Promise;
410 assert.strictEqual(r1a.value, 'chunk1');
411 assert.strictEqual(r2a.value, 'chunk1');
412 
413 // Now cancel branch2
414 await reader2.cancel('done with branch2');
415 
416 // Start another read on branch1
417 read1Promise = reader1.read();
418 
419 // Enqueue second chunk - only branch1 should get it
420 // This should not crash even though branch2's consumer is now closed
421 controller.enqueue('chunk2');
422 const r1b = await read1Promise;
423 assert.strictEqual(r1b.value, 'chunk2');
424 
425 // Start another read and enqueue third chunk to confirm continued operation
426 read1Promise = reader1.read();
427 controller.enqueue('chunk3');
428 const r1c = await read1Promise;
429 assert.strictEqual(r1c.value, 'chunk3');
430 
431 // Close and verify
432 read1Promise = reader1.read();
433 controller.close();
434 const r1d = await read1Promise;
435 assert.strictEqual(r1d.done, true);
436 },
437};
438 
439// ============================================================================
440 
441export const testCancelPipethrough = {
442 async test() {
443 const enc = new TextEncoder();
444 const transform = new IdentityTransformStream();
445 const rs = new ReadableStream({
446 start(c) {
447 c.enqueue(enc.encode('hello'));
448 },
449 });
450 const readable = rs.pipeThrough(transform);
451 
452 const reader = readable.getReader();
453 
454 assert.ok(rs.locked);
455 assert.ok(transform.writable.locked);
456 
457 reader.cancel(new Error('boom'));
458 reader.releaseLock();
459 
460 // We've got to wait a tick to allow the cancel to propagate
461 await scheduler.wait(1);
462 
463 assert.ok(!rs.locked);
464 assert.ok(!transform.readable.locked);
465 assert.ok(!transform.writable.locked);
466 
467 // Our JavaScript ReadableStream should be closed (not errored).
468 // Cancel propagates back and closes the source stream.
469 const reader2 = rs.getReader();
470 const result = await reader2.read();
471 assert.ok(result.done);
472 assert.strictEqual(result.value, undefined);
473 },
474};
475 
476// Same as testCancelPipethrough but uses a JavaScript-backed TransformStream
477// instead of IdentityTransformStream. The behavior should be the same.
478export const testCancelPipethrough2 = {
479 async test() {
480 const enc = new TextEncoder();
481 const transform = new TransformStream();
482 const rs = new ReadableStream({
483 start(c) {
484 c.enqueue(enc.encode('hello'));
485 },
486 });
487 const readable = rs.pipeThrough(transform);
488 
489 const reader = readable.getReader();
490 
491 assert.ok(rs.locked);
492 assert.ok(transform.writable.locked);
493 
494 reader.cancel(new Error('boom'));
495 reader.releaseLock();
496 
497 // We've got to wait a tick to allow the cancel to propagate
498 await scheduler.wait(1);
499 
500 assert.ok(!rs.locked);
501 assert.ok(!transform.readable.locked);
502 assert.ok(!transform.writable.locked);
503 
504 // Our JavaScript ReadableStream should be closed (not errored).
505 // Cancel propagates back and closes the source stream.
506 const reader2 = rs.getReader();
507 const result = await reader2.read();
508 assert.ok(result.done);
509 assert.strictEqual(result.value, undefined);
510 },
511};
512 
513export const ResponseTextLargeBody = {
514 async test() {
515 const targetSize = 2147483648 + 1024 * 1024;
516 let sent = 0;
517 const chunkSize = 64 * 1024 * 1024;
518 
519 const stream = new ReadableStream({
520 pull(controller) {
521 if (sent >= targetSize) {
522 controller.close();
523 return;
524 }
525 const size = Math.min(chunkSize, targetSize - sent);
526 controller.enqueue(new Uint8Array(size).fill(0x41));
527 sent += size;
528 },
529 });
530 
531 await assert.rejects(
532 new Response(stream).text(),
533 (e) => e instanceof RangeError
534 );
535 },
536};