Skip to content
File

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

javascript2614 lines
1// Copyright (c) 2025 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 JavaScript-backed streams (ReadableStream and WritableStream constructors)
6// Ported from edgeworker streams-js.ew-test
7 
8import { strictEqual, ok, throws, rejects } from 'node:assert';
9 
10// Test that JS streams globals exist
11export const userStreamsGlobalsExist = {
12 test() {
13 ok(ReadableStreamDefaultController !== undefined);
14 ok(ReadableByteStreamController !== undefined);
15 ok(ReadableStreamBYOBRequest !== undefined);
16 ok(WritableStreamDefaultController !== undefined);
17 },
18};
19 
20// Test that JS streams objects are not directly constructable
21export const jsStreamsObjectsNotConstructable = {
22 test() {
23 throws(() => new ReadableStreamDefaultController(), TypeError);
24 throws(() => new ReadableByteStreamController(), TypeError);
25 throws(() => new ReadableStreamBYOBRequest(), TypeError);
26 throws(() => new WritableStreamDefaultController(), TypeError);
27 },
28};
29 
30// Test new ReadableStream() works
31export const newReadableStream = {
32 test() {
33 new ReadableStream();
34 new ReadableStream({ type: 'bytes' });
35 },
36};
37 
38// Test that underlying source algorithms are called
39export const newReadableStreamAlgorithms = {
40 async test() {
41 // Sync algorithms
42 {
43 let started = false;
44 let pulled = false;
45 let canceled = false;
46 const rs = new ReadableStream({
47 start() {
48 started = true;
49 },
50 pull() {
51 pulled = true;
52 },
53 cancel() {
54 canceled = true;
55 },
56 });
57 ok(started);
58 
59 await scheduler.wait(1);
60 
61 rs.cancel();
62 
63 ok(pulled);
64 ok(canceled);
65 }
66 
67 // Byte stream sync algorithms
68 {
69 let started = false;
70 let pulled = false;
71 let canceled = false;
72 const rs = new ReadableStream(
73 {
74 type: 'bytes',
75 start() {
76 started = true;
77 },
78 pull() {
79 pulled = true;
80 },
81 cancel() {
82 canceled = true;
83 },
84 },
85 { highWaterMark: 1 }
86 );
87 
88 ok(started);
89 await scheduler.wait(1);
90 
91 rs.cancel();
92 
93 ok(pulled);
94 ok(canceled);
95 }
96 
97 // Async algorithms for value stream
98 {
99 let onStarted, onPulled, onCanceled;
100 let started = new Promise((resolve) => (onStarted = resolve));
101 let pulled = new Promise((resolve) => (onPulled = resolve));
102 let canceled = new Promise((resolve) => (onCanceled = resolve));
103 
104 const rs = new ReadableStream({
105 async start() {
106 await scheduler.wait(1);
107 onStarted();
108 },
109 async pull() {
110 await scheduler.wait(1);
111 onPulled();
112 },
113 async cancel() {
114 onCanceled();
115 },
116 });
117 
118 await Promise.allSettled([started, pulled]);
119 await scheduler.wait(1);
120 await Promise.allSettled([rs.cancel(), canceled]);
121 }
122 
123 // Async algorithms for byte stream
124 {
125 let onStarted, onPulled, onCanceled;
126 let started = new Promise((resolve) => (onStarted = resolve));
127 let pulled = new Promise((resolve) => (onPulled = resolve));
128 let canceled = new Promise((resolve) => (onCanceled = resolve));
129 
130 const rs = new ReadableStream(
131 {
132 type: 'bytes',
133 async start() {
134 await scheduler.wait(1);
135 onStarted();
136 },
137 async pull() {
138 await scheduler.wait(1);
139 onPulled();
140 },
141 async cancel() {
142 onCanceled();
143 },
144 },
145 { highWaterMark: 1 }
146 );
147 
148 await Promise.allSettled([started, pulled]);
149 await scheduler.wait(1);
150 await Promise.allSettled([rs.cancel(), canceled]);
151 }
152 },
153};
154 
155// Test that new ReadableStream creates the right kind of controller
156export const newReadableStreamControllerType = {
157 test() {
158 new ReadableStream({
159 start(c) {
160 ok(c instanceof ReadableStreamDefaultController);
161 },
162 pull(c) {
163 ok(c instanceof ReadableStreamDefaultController);
164 },
165 });
166 
167 new ReadableStream({
168 type: 'bytes',
169 start(c) {
170 ok(c instanceof ReadableByteStreamController);
171 },
172 pull(c) {
173 ok(c instanceof ReadableByteStreamController);
174 const byobRequest = c.byobRequest;
175 ok(byobRequest != null);
176 ok(byobRequest === c.byobRequest);
177 ok(byobRequest instanceof ReadableStreamBYOBRequest);
178 ok(byobRequest.view instanceof Uint8Array);
179 },
180 });
181 },
182};
183 
184// Test sync algorithm errors are handled properly
185export const newReadableStreamSyncAlgorithmErrorsHandled = {
186 async test() {
187 // Start error
188 {
189 const rs = new ReadableStream({
190 start() {
191 throw new Error('boom');
192 },
193 });
194 
195 await rejects(rs.getReader().read(), { message: 'boom' });
196 }
197 
198 // Pull error
199 {
200 const rs = new ReadableStream({
201 pull() {
202 throw new Error('boom');
203 },
204 });
205 
206 await rejects(rs.getReader().read(), { message: 'boom' });
207 }
208 
209 // Cancel error
210 {
211 const rs = new ReadableStream({
212 cancel() {
213 throw new Error('boom');
214 },
215 });
216 await rejects(rs.cancel(), { message: 'boom' });
217 }
218 },
219};
220 
221// Test async algorithm errors are handled properly
222export const newReadableStreamAsyncAlgorithmErrorsHandled = {
223 async test() {
224 // Async start error
225 {
226 const rs = new ReadableStream({
227 async start() {
228 throw new Error('boom');
229 },
230 });
231 
232 await rejects(rs.getReader().read(), { message: 'boom' });
233 }
234 
235 // Async pull error
236 {
237 const rs = new ReadableStream({
238 async pull() {
239 throw new Error('boom');
240 },
241 });
242 
243 await rejects(rs.getReader().read(), { message: 'boom' });
244 }
245 
246 // Async cancel error
247 {
248 const rs = new ReadableStream({
249 async cancel() {
250 throw new Error('boom');
251 },
252 });
253 
254 await rejects(rs.cancel(), { message: 'boom' });
255 }
256 },
257};
258 
259// Test size algorithm is called with correct value and errors handled
260export const sizeAlgorithmCalled = {
261 async test() {
262 // Size algorithm called with correct value
263 {
264 let sizeCalled = false;
265 new ReadableStream(
266 {
267 pull(c) {
268 c.enqueue(1);
269 },
270 },
271 {
272 size(value) {
273 strictEqual(value, 1);
274 sizeCalled = true;
275 },
276 }
277 );
278 
279 ok(sizeCalled);
280 }
281 
282 // Size algorithm ignored in byte streams
283 {
284 let sizeCalled = false;
285 new ReadableStream(
286 {
287 type: 'bytes',
288 pull(c) {
289 c.enqueue(new Uint8Array(1));
290 },
291 },
292 {
293 size() {
294 sizeCalled = true;
295 },
296 }
297 );
298 
299 ok(!sizeCalled);
300 }
301 
302 // Size algorithm error handled
303 {
304 const rs = new ReadableStream(
305 {
306 pull(c) {
307 c.enqueue(1);
308 },
309 },
310 {
311 size() {
312 throw new Error('boom');
313 },
314 }
315 );
316 
317 await rejects(rs.getReader().read(), { message: 'boom' });
318 }
319 
320 // Async size algorithm not allowed
321 {
322 const rs = new ReadableStream(
323 {
324 pull(c) {
325 c.enqueue(1);
326 },
327 },
328 {
329 async size() {
330 return 1;
331 },
332 }
333 );
334 
335 await rejects(rs.getReader().read(), {
336 message: 'The value cannot be converted because it is not an integer.',
337 });
338 }
339 },
340};
341 
342// Test ReadableStream getDesiredSize is calculated correctly
343export const readableGetDesiredSize = {
344 async test() {
345 // Value stream desiredSize
346 {
347 let controller;
348 
349 const rs = new ReadableStream(
350 {
351 start(c) {
352 controller = c;
353 strictEqual(c.desiredSize, 2);
354 c.enqueue(1);
355 strictEqual(c.desiredSize, 1);
356 c.enqueue(2);
357 strictEqual(c.desiredSize, 0);
358 c.enqueue(3);
359 strictEqual(c.desiredSize, -1);
360 },
361 },
362 {
363 highWaterMark: 2,
364 }
365 );
366 
367 await rs.getReader().read();
368 strictEqual(controller.desiredSize, 0);
369 }
370 
371 // Enqueuing when there's an active read skips the queue
372 {
373 let controller;
374 const rs = new ReadableStream(
375 {
376 start(c) {
377 controller = c;
378 },
379 },
380 { highWaterMark: 2 }
381 );
382 
383 const reader = rs.getReader();
384 strictEqual(controller.desiredSize, 2);
385 const read = reader.read();
386 controller.enqueue(1);
387 strictEqual(controller.desiredSize, 2);
388 strictEqual((await read).value, 1);
389 }
390 
391 // Byte stream desiredSize
392 {
393 let controller;
394 const rs = new ReadableStream(
395 {
396 type: 'bytes',
397 start(c) {
398 controller = c;
399 strictEqual(c.desiredSize, 2);
400 c.enqueue(new Uint8Array(2));
401 strictEqual(c.desiredSize, 0);
402 c.enqueue(new Uint8Array(1));
403 strictEqual(c.desiredSize, -1);
404 },
405 },
406 {
407 highWaterMark: 2,
408 }
409 );
410 
411 strictEqual((await rs.getReader().read()).value.byteLength, 3);
412 strictEqual(controller.desiredSize, 2);
413 }
414 
415 // Byte stream enqueuing when there's an active read skips the queue
416 {
417 let controller;
418 const rs = new ReadableStream(
419 {
420 type: 'bytes',
421 start(c) {
422 controller = c;
423 },
424 },
425 { highWaterMark: 2 }
426 );
427 
428 const reader = rs.getReader();
429 strictEqual(controller.desiredSize, 2);
430 const read = reader.read();
431 controller.enqueue(new Uint8Array(10));
432 strictEqual(controller.desiredSize, 2);
433 strictEqual((await read).value.byteLength, 10);
434 }
435 },
436};
437 
438// Test ReadableStream controller.error() works as expected
439export const readableStreamControllerError = {
440 async test() {
441 // Value stream
442 {
443 let controller;
444 const rs = new ReadableStream({
445 start(c) {
446 controller = c;
447 },
448 });
449 const reader = rs.getReader();
450 const read = reader.read();
451 controller.error(new Error('bang!'));
452 await rejects(read, { message: 'bang!' });
453 }
454 
455 // Byte stream
456 {
457 let controller;
458 const rs = new ReadableStream({
459 type: 'bytes',
460 start(c) {
461 controller = c;
462 },
463 });
464 const reader = rs.getReader();
465 const read = reader.read();
466 controller.error(new Error('bang!'));
467 await rejects(read, { message: 'bang!' });
468 }
469 },
470};
471 
472// Test ReadableStream autoAllocateChunkSize works as expected
473export const readableStreamAutoAllocateChunkSize = {
474 async test() {
475 throws(() => {
476 new ReadableStream({
477 type: 'bytes',
478 autoAllocateChunkSize: 0,
479 });
480 }, TypeError);
481 
482 throws(() => {
483 new ReadableStream({
484 type: 'bytes',
485 autoAllocateChunkSize: -1,
486 });
487 }, TypeError);
488 
489 throws(() => {
490 new ReadableStream({
491 type: 'bytes',
492 autoAllocateChunkSize: 'a',
493 });
494 }, TypeError);
495 
496 let pulled = false;
497 const rs = new ReadableStream({
498 type: 'bytes',
499 autoAllocateChunkSize: 10,
500 pull(c) {
501 pulled = true;
502 if (c.byobRequest) {
503 strictEqual(c.byobRequest.view.byteLength, 10);
504 c.byobRequest.respond(10);
505 }
506 },
507 });
508 await rs.getReader().read();
509 ok(pulled);
510 },
511};
512 
513// Test ReadableStream byte stream respond() works appropriately
514export const readableStreamByteRespond = {
515 async test() {
516 // Basic respond
517 {
518 const rs = new ReadableStream({
519 type: 'bytes',
520 pull(c) {
521 if (c.byobRequest) {
522 const req = c.byobRequest;
523 req.view[0] = 1;
524 req.view[1] = 2;
525 req.view[2] = 3;
526 
527 throws(() => req.respond(10), RangeError);
528 throws(() => req.respond(0), TypeError);
529 
530 req.respond(3);
531 
532 // This will error the stream but won't be immediately
533 // apparent until the next read operation.
534 req.respond(3);
535 }
536 },
537 });
538 
539 const reader = rs.getReader({ mode: 'byob' });
540 const u8 = new Uint8Array(3);
541 const read = reader.read(u8);
542 strictEqual(u8.byteLength, 0);
543 
544 const { value } = await read;
545 strictEqual(value.byteLength, 3);
546 
547 await rejects(reader.read(new Uint8Array(3)), {
548 message: 'This ReadableStreamBYOBRequest has been invalidated.',
549 });
550 }
551 
552 // Respond with close
553 {
554 const rs = new ReadableStream({
555 type: 'bytes',
556 pull(c) {
557 if (c.byobRequest) {
558 c.close();
559 c.byobRequest.respond(0);
560 }
561 },
562 });
563 
564 const reader = rs.getReader({ mode: 'byob' });
565 
566 const u8 = new Uint8Array([1, 2, 3]);
567 
568 const { done, value } = await reader.read(u8);
569 
570 ok(done);
571 ok(value instanceof Uint8Array);
572 strictEqual(value.byteLength, 0);
573 strictEqual(value.buffer.byteLength, 3);
574 const u82 = new Uint8Array(value.buffer, 0, 3);
575 strictEqual(u82[0], 1);
576 strictEqual(u82[1], 2);
577 strictEqual(u82[2], 3);
578 }
579 },
580};
581 
582// Test ReadableStream byte stream respondWithNewView works appropriately
583export const readableStreamByteRespondWithNewView = {
584 async test() {
585 // Basic respondWithNewView
586 {
587 const rs = new ReadableStream({
588 type: 'bytes',
589 pull(c) {
590 if (c.byobRequest) {
591 const req = c.byobRequest;
592 const u8 = new Uint8Array(req.view.buffer);
593 
594 u8[0] = 1;
595 u8[1] = 2;
596 u8[2] = 3;
597 
598 // Can't respond with zero if we're not closed.
599 throws(() => req.respondWithNewView(new Uint8Array(0)), TypeError);
600 
601 // Underlying buffer is too big.
602 throws(
603 () => req.respondWithNewView(new Uint8Array(10)),
604 RangeError
605 );
606 
607 // Can't respond with a non-detachable ArrayBuffer.
608 throws(
609 () =>
610 req.respondWithNewView(
611 new Uint8Array(new SharedArrayBuffer(10))
612 ),
613 TypeError
614 );
615 
616 // New view has an invalid byte offset.
617 throws(
618 () => req.respondWithNewView(new Uint8Array(req.view.buffer, 1)),
619 RangeError
620 );
621 
622 req.respondWithNewView(u8);
623 
624 strictEqual(u8.byteLength, 0);
625 
626 // This will error the stream but won't be immediately
627 // apparent until the next read operation.
628 req.respond(3);
629 }
630 },
631 });
632 
633 const reader = rs.getReader({ mode: 'byob' });
634 const u8 = new Uint8Array(3);
635 const read = reader.read(u8);
636 strictEqual(u8.byteLength, 0);
637 
638 const { value } = await read;
639 strictEqual(value.byteLength, 3);
640 strictEqual(value[0], 1);
641 strictEqual(value[1], 2);
642 strictEqual(value[2], 3);
643 
644 await rejects(reader.read(new Uint8Array(3)), {
645 message: 'This ReadableStreamBYOBRequest has been invalidated.',
646 });
647 }
648 
649 // RespondWithNewView with close
650 {
651 const rs = new ReadableStream({
652 type: 'bytes',
653 pull(c) {
654 if (c.byobRequest) {
655 c.close();
656 c.byobRequest.respondWithNewView(
657 new Uint8Array(c.byobRequest.view.buffer, 0, 0)
658 );
659 }
660 },
661 });
662 
663 const reader = rs.getReader({ mode: 'byob' });
664 
665 const { done, value } = await reader.read(new Uint8Array(3));
666 
667 ok(done);
668 ok(value instanceof Uint8Array);
669 strictEqual(value.byteLength, 0);
670 strictEqual(value.buffer.byteLength, 3);
671 }
672 },
673};
674 
675// Test ReadableStream JS controllers allow for multiple pending reads
676export const readableStreamMultiplePendingReads = {
677 async test() {
678 // Value stream
679 {
680 let controller;
681 const rs = new ReadableStream({
682 start(c) {
683 controller = c;
684 },
685 });
686 const reader = rs.getReader();
687 const read1 = reader.read();
688 const read2 = reader.read();
689 controller.enqueue(1);
690 controller.enqueue(2);
691 const [res1, res2] = await Promise.all([read1, read2]);
692 strictEqual(res1.value, 1);
693 strictEqual(res2.value, 2);
694 }
695 
696 // Byte stream
697 {
698 let controller;
699 const rs = new ReadableStream({
700 type: 'bytes',
701 start(c) {
702 controller = c;
703 },
704 });
705 const enc = new TextEncoder();
706 const dec = new TextDecoder();
707 const reader = rs.getReader();
708 const read1 = reader.read();
709 const read2 = reader.read();
710 controller.enqueue(enc.encode('hello'));
711 controller.enqueue(enc.encode('there'));
712 const [res1, res2] = await Promise.all([read1, read2]);
713 strictEqual(dec.decode(res1.value), 'hello');
714 strictEqual(dec.decode(res2.value), 'there');
715 }
716 },
717};
718 
719// Test ReadableStream byte controller enqueue and reads with mismatched sizes works
720export const readableStreamBytesMismatchedSizes = {
721 async test() {
722 const enc = new TextEncoder();
723 let pulls = 0;
724 const rs = new ReadableStream({
725 type: 'bytes',
726 start(c) {
727 c.enqueue(enc.encode('hello'));
728 },
729 pull(c) {
730 if (c.byobRequest) {
731 pulls++;
732 c.enqueue(enc.encode('there'));
733 }
734 },
735 });
736 const reader = rs.getReader({ mode: 'byob' });
737 
738 await Promise.all(
739 [
740 enc.encode('he'),
741 enc.encode('ll'),
742 enc.encode('o'),
743 enc.encode('th'),
744 enc.encode('er'),
745 enc.encode('e'),
746 ].map(async (i) => {
747 const { done, value } = await reader.read(new Uint8Array(2));
748 ok(!done);
749 strictEqual(value.byteLength, i.byteLength);
750 for (let n = 0; n < value.byteLength; n++) {
751 strictEqual(value[n], i[n]);
752 }
753 })
754 );
755 
756 strictEqual(pulls, 1);
757 },
758};
759 
760// Test ReadableStream byte controller enqueue and reads with mismatched view types works
761export const readableStreamBytesMismatchedViewTypes = {
762 async test() {
763 let pull = 0;
764 const rs = new ReadableStream({
765 type: 'bytes',
766 pull(c) {
767 if (c.byobRequest) {
768 const view = c.byobRequest.view;
769 switch (pull++) {
770 case 0: {
771 strictEqual(view.byteLength, 8);
772 strictEqual(view.byteOffset, 0);
773 view[0] = 1;
774 view[1] = 2;
775 view[2] = 3;
776 view[3] = 4;
777 view[4] = 5;
778 view[5] = 6;
779 view[6] = 7;
780 c.byobRequest.respond(7);
781 break;
782 }
783 case 1: {
784 strictEqual(view.byteLength, 5);
785 strictEqual(view.byteOffset, 3);
786 view[0] = 8;
787 c.byobRequest.respond(1);
788 c.close();
789 break;
790 }
791 }
792 }
793 },
794 });
795 
796 const r = rs.getReader({ mode: 'byob' });
797 
798 {
799 const { value } = await r.read(new Uint32Array(2));
800 ok(value instanceof Uint32Array);
801 strictEqual(value.length, 1);
802 strictEqual(value.byteLength, 4);
803 strictEqual(value.buffer.byteLength, 8);
804 const u8 = new Uint8Array(value.buffer, 0, value.byteLength);
805 strictEqual(u8[0], 1);
806 strictEqual(u8[1], 2);
807 strictEqual(u8[2], 3);
808 strictEqual(u8[3], 4);
809 }
810 
811 {
812 const { value } = await r.read(new Uint32Array(2));
813 ok(value instanceof Uint32Array);
814 strictEqual(value.length, 1);
815 strictEqual(value.byteLength, 4);
816 strictEqual(value.buffer.byteLength, 8);
817 const u8 = new Uint8Array(value.buffer);
818 strictEqual(u8[0], 5);
819 strictEqual(u8[1], 6);
820 strictEqual(u8[2], 7);
821 strictEqual(u8[3], 8);
822 }
823 },
824};
825 
826// Test ReadableStream byte controller enqueue subarray works
827export const readableStreamBytesEnqueueSubarray = {
828 async test() {
829 const enc = new TextEncoder();
830 const dec = new TextDecoder();
831 const rs = new ReadableStream({
832 type: 'bytes',
833 pull(c) {
834 const u8 = enc.encode('hello');
835 c.enqueue(u8.subarray(1, 4));
836 strictEqual(u8.byteLength, 0);
837 c.close();
838 },
839 });
840 
841 const r = rs.getReader({ mode: 'byob' });
842 
843 const { value } = await r.read(new Uint8Array(5));
844 
845 strictEqual(dec.decode(value), 'ell');
846 },
847};
848 
849// Test ReadableStream default and bytes controllers close promise works
850export const readableStreamDefaultClosePromise = {
851 async test() {
852 // Value stream
853 {
854 let controller;
855 const rs = new ReadableStream({
856 start(c) {
857 controller = c;
858 },
859 });
860 const r = rs.getReader();
861 let closed = false;
862 r.closed.then(() => (closed = true));
863 controller.enqueue(1);
864 controller.close();
865 await r.read();
866 ok(closed);
867 }
868 
869 // Byte stream default reader
870 {
871 let controller;
872 const rs = new ReadableStream({
873 type: 'bytes',
874 start(c) {
875 controller = c;
876 },
877 });
878 
879 const r = rs.getReader();
880 
881 let closed = false;
882 r.closed.then(() => (closed = true));
883 controller.enqueue(new Uint8Array(1));
884 controller.close();
885 await r.read();
886 await scheduler.wait(1);
887 ok(closed);
888 }
889 
890 // Byte stream BYOB reader
891 {
892 let controller;
893 const rs = new ReadableStream({
894 type: 'bytes',
895 start(c) {
896 controller = c;
897 },
898 });
899 
900 const r = rs.getReader({ mode: 'byob' });
901 
902 let closed = false;
903 r.closed.then(() => (closed = true));
904 controller.enqueue(new Uint8Array(1));
905 controller.close();
906 await r.read(new Uint8Array(1));
907 await scheduler.wait(1);
908 ok(closed);
909 }
910 },
911};
912 
913// Test ReadableStream default and bytes reads can be canceled
914export const readableStreamCancelReads = {
915 async test() {
916 // Value stream
917 {
918 const rs = new ReadableStream();
919 const reader = rs.getReader();
920 const read = reader.read();
921 reader.cancel();
922 
923 const { done, value } = await read;
924 ok(done);
925 strictEqual(value, undefined);
926 }
927 
928 // Byte stream default reader
929 {
930 const rs = new ReadableStream({
931 type: 'bytes',
932 });
933 const reader = rs.getReader();
934 const read = reader.read();
935 reader.cancel();
936 
937 const { done } = await read;
938 ok(done);
939 }
940 
941 // Byte stream BYOB reader
942 {
943 const rs = new ReadableStream({
944 type: 'bytes',
945 });
946 const reader = rs.getReader({ mode: 'byob' });
947 const read = reader.read(new Uint8Array(1));
948 reader.cancel();
949 
950 const { done } = await read;
951 ok(done);
952 }
953 
954 // Byte stream BYOB reader with cancel reason
955 {
956 let cancelCalled = false;
957 const rs = new ReadableStream({
958 type: 'bytes',
959 cancel(reason) {
960 strictEqual(reason, 'boom');
961 cancelCalled = true;
962 },
963 });
964 const reader = rs.getReader({ mode: 'byob' });
965 const read = reader.read(new Uint8Array(1));
966 reader.cancel('boom');
967 
968 const { done } = await read;
969 ok(done);
970 ok(cancelCalled);
971 }
972 },
973};
974 
975// Test ReadableStream default and byte controller release lock work
976export const readableStreamReleaseLock = {
977 async test() {
978 // With capture_async_api_throws, async methods (pipeTo) return rejected promises instead of throwing
979 // pipeThrough returns ReadableStream (not a promise), so it still throws synchronously
980 const captureAsyncThrows =
981 Cloudflare.compatibilityFlags.capture_async_api_throws;
982 
983 // Value stream
984 {
985 const rs = new ReadableStream();
986 const reader = rs.getReader();
987 throws(() => rs.getReader(), TypeError);
988 throws(() => rs.tee(), TypeError);
989 if (captureAsyncThrows) {
990 await rejects(rs.pipeTo(), TypeError);
991 } else {
992 throws(() => rs.pipeTo(), TypeError);
993 }
994 throws(() => rs.pipeThrough(), TypeError);
995 
996 reader.releaseLock();
997 rs.getReader();
998 }
999 
1000 // Byte stream default reader
1001 {
1002 const rs = new ReadableStream({
1003 type: 'bytes',
1004 });
1005 const reader = rs.getReader();
1006 throws(() => rs.getReader(), TypeError);
1007 throws(() => rs.tee(), TypeError);
1008 if (captureAsyncThrows) {
1009 await rejects(rs.pipeTo(), TypeError);
1010 } else {
1011 throws(() => rs.pipeTo(), TypeError);
1012 }
1013 throws(() => rs.pipeThrough(), TypeError);
1014 
1015 reader.releaseLock();
1016 rs.getReader();
1017 }
1018 
1019 // Byte stream BYOB reader
1020 {
1021 const rs = new ReadableStream({
1022 type: 'bytes',
1023 });
1024 const reader = rs.getReader({ mode: 'byob' });
1025 throws(() => rs.getReader(), TypeError);
1026 throws(() => rs.tee(), TypeError);
1027 if (captureAsyncThrows) {
1028 await rejects(rs.pipeTo(), TypeError);
1029 } else {
1030 throws(() => rs.pipeTo(), TypeError);
1031 }
1032 throws(() => rs.pipeThrough(), TypeError);
1033 
1034 reader.releaseLock();
1035 rs.getReader();
1036 }
1037 },
1038};
1039 
1040// Test ReadableStream default controller does not support BYOB reader
1041export const readableStreamDefaultNoByob = {
1042 test() {
1043 const rs = new ReadableStream();
1044 throws(() => rs.getReader({ mode: 'byob' }), TypeError);
1045 throws(() => new ReadableStreamBYOBReader(rs), TypeError);
1046 },
1047};
1048 
1049// Test ReadableStream default controller tee() works
1050export const readableStreamDefaultTee = {
1051 async test() {
1052 // Tee an immediately closed ReadableStream
1053 {
1054 const rs = new ReadableStream({
1055 start(c) {
1056 c.close();
1057 },
1058 });
1059 
1060 const [branch1, branch2] = rs.tee();
1061 
1062 const reader1 = branch1.getReader();
1063 const reader2 = branch2.getReader();
1064 
1065 const [res1, res2] = await Promise.all([reader1.read(), reader2.read()]);
1066 
1067 strictEqual(res1.done, true);
1068 strictEqual(res2.done, true);
1069 }
1070 
1071 // Tee with data
1072 {
1073 const rs = new ReadableStream({
1074 pull(c) {
1075 c.enqueue(1);
1076 c.close();
1077 },
1078 });
1079 
1080 const [branch1, branch2] = rs.tee();
1081 
1082 const reader1 = branch1.getReader();
1083 const reader2 = branch2.getReader();
1084 
1085 const [res1, res2] = await Promise.all([reader1.read(), reader2.read()]);
1086 
1087 strictEqual(res1.value, 1);
1088 strictEqual(res2.value, 1);
1089 
1090 const [res3, res4] = await Promise.all([reader1.read(), reader2.read()]);
1091 
1092 strictEqual(res3.done, true);
1093 strictEqual(res4.done, true);
1094 }
1095 
1096 // Tee with multiple enqueues
1097 {
1098 let counter = 0;
1099 const rs = new ReadableStream({
1100 pull(c) {
1101 c.enqueue(counter++);
1102 if (counter == 2) {
1103 c.close();
1104 }
1105 },
1106 });
1107 
1108 const [branch1, branch2] = rs.tee();
1109 
1110 const reader1 = branch1.getReader();
1111 const reader2 = branch2.getReader();
1112 
1113 {
1114 const [result1, result2] = await Promise.all([
1115 reader1.read(),
1116 reader2.read(),
1117 ]);
1118 
1119 ok(!result1.done);
1120 ok(!result2.done);
1121 strictEqual(result1.value, 0);
1122 strictEqual(result2.value, 0);
1123 }
1124 
1125 {
1126 const [result1, result2] = await Promise.all([
1127 reader1.read(),
1128 reader2.read(),
1129 ]);
1130 ok(!result1.done);
1131 ok(!result2.done);
1132 strictEqual(result1.value, 1);
1133 strictEqual(result2.value, 1);
1134 }
1135 
1136 {
1137 const [result1, result2] = await Promise.all([
1138 reader1.read(),
1139 reader2.read(),
1140 ]);
1141 ok(result1.done);
1142 ok(result2.done);
1143 strictEqual(result1.value, undefined);
1144 strictEqual(result2.value, undefined);
1145 }
1146 }
1147 
1148 // Canceling one branch does not impact the other
1149 {
1150 let counter = 0;
1151 let canceled = false;
1152 const rs = new ReadableStream({
1153 pull(c) {
1154 c.enqueue(counter++);
1155 if (counter == 2) {
1156 c.close();
1157 }
1158 },
1159 cancel() {
1160 canceled = true;
1161 },
1162 });
1163 
1164 const [branch1, branch2] = rs.tee();
1165 
1166 const reader1 = branch1.getReader();
1167 const reader2 = branch2.getReader();
1168 
1169 {
1170 const [result1, result2] = await Promise.all([
1171 reader1.read(),
1172 reader2.read(),
1173 ]);
1174 
1175 ok(!result1.done);
1176 ok(!result2.done);
1177 strictEqual(result1.value, 0);
1178 strictEqual(result2.value, 0);
1179 }
1180 
1181 reader2.cancel();
1182 
1183 {
1184 const [result1, result2] = await Promise.all([
1185 reader1.read(),
1186 reader2.read(),
1187 ]);
1188 
1189 ok(!canceled);
1190 
1191 ok(!result1.done);
1192 ok(result2.done);
1193 strictEqual(result1.value, 1);
1194 strictEqual(result2.value, undefined);
1195 }
1196 
1197 {
1198 const [result1, result2] = await Promise.all([
1199 reader1.read(),
1200 reader2.read(),
1201 ]);
1202 
1203 ok(result1.done);
1204 ok(result2.done);
1205 strictEqual(result1.value, undefined);
1206 strictEqual(result2.value, undefined);
1207 }
1208 }
1209 
1210 // Canceling both tee branches cancels the underlying source
1211 {
1212 let canceled = false;
1213 const rs = new ReadableStream({
1214 start(c) {
1215 c.enqueue(0);
1216 },
1217 cancel() {
1218 canceled = true;
1219 },
1220 });
1221 
1222 const [branch1, branch2] = rs.tee();
1223 
1224 const reader1 = branch1.getReader();
1225 const reader2 = branch2.getReader();
1226 
1227 {
1228 const [result1, result2] = await Promise.all([
1229 reader1.read(),
1230 reader2.read(),
1231 ]);
1232 
1233 ok(!result1.done);
1234 ok(!result2.done);
1235 strictEqual(result1.value, 0);
1236 strictEqual(result2.value, 0);
1237 }
1238 
1239 await reader1.cancel();
1240 ok(!canceled);
1241 
1242 await reader2.cancel();
1243 ok(canceled);
1244 }
1245 
1246 // Tee of a tee works
1247 {
1248 let controller;
1249 const rs = new ReadableStream({
1250 start(c) {
1251 controller = c;
1252 },
1253 });
1254 
1255 const [branch1, branch2] = rs.tee();
1256 const [branch3, branch4] = branch2.tee();
1257 
1258 throws(() => branch2.getReader(), TypeError);
1259 
1260 {
1261 const reader1 = branch1.getReader();
1262 const reader3 = branch3.getReader();
1263 const reader4 = branch4.getReader();
1264 
1265 const read1 = reader1.read();
1266 const read3 = reader3.read();
1267 const read4 = reader4.read();
1268 
1269 controller.enqueue(1);
1270 
1271 const [result1, result3, result4] = await Promise.all([
1272 read1,
1273 read3,
1274 read4,
1275 ]);
1276 
1277 strictEqual(result1.value, 1);
1278 strictEqual(result3.value, 1);
1279 strictEqual(result4.value, 1);
1280 }
1281 }
1282 
1283 // Erroring the underlying source errors the branches
1284 {
1285 let controller;
1286 const rs = new ReadableStream({
1287 start(c) {
1288 controller = c;
1289 },
1290 });
1291 
1292 const [branch1, branch2] = rs.tee();
1293 
1294 const reader1 = branch1.getReader();
1295 const reader2 = branch2.getReader();
1296 
1297 const read1 = reader1.read();
1298 const read2 = reader2.read();
1299 
1300 controller.error('boom');
1301 
1302 (await Promise.allSettled([read1, read2])).forEach((i) => {
1303 strictEqual(i.status, 'rejected');
1304 strictEqual(i.reason, 'boom');
1305 });
1306 }
1307 
1308 // Tee branches support BYOB reads
1309 {
1310 let controller;
1311 const enc = new TextEncoder();
1312 const dec = new TextDecoder();
1313 const rs = new ReadableStream({
1314 type: 'bytes',
1315 start(c) {
1316 controller = c;
1317 },
1318 });
1319 
1320 const [branch1, branch2] = rs.tee();
1321 
1322 const reader1 = branch1.getReader({ mode: 'byob' });
1323 const reader2 = branch2.getReader({ mode: 'byob' });
1324 
1325 const buf1 = new Uint8Array(2);
1326 const buf2 = new Uint8Array(3);
1327 
1328 const promises = [reader1.read(buf1), reader2.read(buf2)];
1329 
1330 controller.enqueue(enc.encode('hello'));
1331 
1332 const results = await Promise.all(promises);
1333 
1334 strictEqual(dec.decode(results[0].value), 'he');
1335 strictEqual(dec.decode(results[1].value), 'hel');
1336 }
1337 },
1338};
1339 
1340// =====================================================================================
1341// WritableStream tests
1342// =====================================================================================
1343 
1344// Test new WritableStream() works
1345export const newWritableStream = {
1346 test() {
1347 new WritableStream();
1348 },
1349};
1350 
1351// Test new WritableStream() with sink works
1352export const newWritableStreamWithSink = {
1353 async test() {
1354 // Sync sink with abort
1355 {
1356 let started = false;
1357 let written = false;
1358 let closed = false;
1359 let aborted = false;
1360 const ws = new WritableStream({
1361 start(c) {
1362 ok(c instanceof WritableStreamDefaultController);
1363 started = true;
1364 },
1365 write(value, c) {
1366 strictEqual(value, 1);
1367 ok(c instanceof WritableStreamDefaultController);
1368 written = true;
1369 },
1370 abort(reason) {
1371 strictEqual(reason.message, 'boom');
1372 aborted = true;
1373 },
1374 close() {
1375 closed = true;
1376 },
1377 });
1378 
1379 ok(started);
1380 
1381 const writer = ws.getWriter();
1382 
1383 await writer.write(1);
1384 ok(written);
1385 
1386 await writer.abort(new Error('boom'));
1387 ok(aborted);
1388 ok(!closed);
1389 
1390 await rejects(writer.closed);
1391 }
1392 
1393 // Sync sink with close
1394 {
1395 let started = false;
1396 let written = false;
1397 let closed = false;
1398 let aborted = false;
1399 const ws = new WritableStream({
1400 start(c) {
1401 ok(c instanceof WritableStreamDefaultController);
1402 started = true;
1403 },
1404 write(value, c) {
1405 strictEqual(value, 1);
1406 ok(c instanceof WritableStreamDefaultController);
1407 written = true;
1408 },
1409 abort() {
1410 aborted = true;
1411 },
1412 close() {
1413 closed = true;
1414 },
1415 });
1416 
1417 ok(started);
1418 
1419 const writer = ws.getWriter();
1420 
1421 await writer.write(1);
1422 ok(written);
1423 
1424 await Promise.all([writer.close(), writer.closed]);
1425 
1426 ok(!aborted);
1427 ok(closed);
1428 }
1429 },
1430};
1431 
1432// Test new WritableStream() with async sink works
1433export const newWritableStreamWithSinkAsync = {
1434 async test() {
1435 // Async sink with abort
1436 {
1437 let started = false;
1438 let written = false;
1439 let closed = false;
1440 let aborted = false;
1441 const ws = new WritableStream({
1442 async start(c) {
1443 ok(c instanceof WritableStreamDefaultController);
1444 await scheduler.wait(10);
1445 started = true;
1446 },
1447 async write(value, c) {
1448 await scheduler.wait(10);
1449 strictEqual(value, 1);
1450 ok(c instanceof WritableStreamDefaultController);
1451 written = true;
1452 },
1453 async abort(reason) {
1454 await scheduler.wait(10);
1455 strictEqual(reason.message, 'boom');
1456 aborted = true;
1457 },
1458 async close() {
1459 closed = true;
1460 },
1461 });
1462 
1463 await scheduler.wait(15);
1464 ok(started);
1465 
1466 const writer = ws.getWriter();
1467 await writer.ready;
1468 
1469 await writer.write(1);
1470 ok(written);
1471 
1472 await writer.abort(new Error('boom'));
1473 ok(aborted);
1474 ok(!closed);
1475 
1476 await rejects(writer.closed);
1477 }
1478 
1479 // Async sink with close
1480 {
1481 let started = false;
1482 let written = false;
1483 let closed = false;
1484 let aborted = false;
1485 const ws = new WritableStream({
1486 async start(c) {
1487 ok(c instanceof WritableStreamDefaultController);
1488 await scheduler.wait(10);
1489 started = true;
1490 },
1491 async write(value, c) {
1492 await scheduler.wait(10);
1493 strictEqual(value, 1);
1494 ok(c instanceof WritableStreamDefaultController);
1495 written = true;
1496 },
1497 async abort() {
1498 await scheduler.wait(10);
1499 aborted = true;
1500 },
1501 async close() {
1502 await scheduler.wait(10);
1503 closed = true;
1504 },
1505 });
1506 
1507 await scheduler.wait(15);
1508 ok(started);
1509 
1510 const writer = ws.getWriter();
1511 await writer.ready;
1512 
1513 await writer.write(1);
1514 ok(written);
1515 
1516 await Promise.all([writer.close(), writer.closed]);
1517 ok(!aborted);
1518 ok(closed);
1519 }
1520 },
1521};
1522 
1523// Test new WritableStream() start algorithm error handled
1524export const newWritableStreamStartError = {
1525 async test() {
1526 // Sync start error
1527 {
1528 const ws = new WritableStream({
1529 start() {
1530 throw new Error('boom');
1531 },
1532 });
1533 
1534 const writer = ws.getWriter();
1535 await rejects(writer.write(1), { message: 'boom' });
1536 }
1537 
1538 // Async start error
1539 {
1540 const ws = new WritableStream({
1541 async start() {
1542 throw new Error('boom');
1543 },
1544 });
1545 
1546 const writer = ws.getWriter();
1547 await rejects(writer.write(1), { message: 'boom' });
1548 }
1549 
1550 // Start with controller.error
1551 {
1552 const ws = new WritableStream({
1553 start(c) {
1554 c.error(new Error('boom'));
1555 },
1556 });
1557 
1558 const writer = ws.getWriter();
1559 await rejects(writer.write(1), { message: 'boom' });
1560 }
1561 },
1562};
1563 
1564// Test new WritableStream() write algorithm error handled
1565export const newWritableStreamWriteError = {
1566 async test() {
1567 // Sync write error
1568 {
1569 const ws = new WritableStream({
1570 write() {
1571 throw new Error('boom');
1572 },
1573 });
1574 
1575 const writer = ws.getWriter();
1576 await rejects(writer.write(1), { message: 'boom' });
1577 }
1578 
1579 // Async write error
1580 {
1581 const ws = new WritableStream({
1582 async write() {
1583 throw new Error('boom');
1584 },
1585 });
1586 
1587 const writer = ws.getWriter();
1588 await rejects(writer.write(1), { message: 'boom' });
1589 }
1590 
1591 // Write with controller.error
1592 {
1593 const ws = new WritableStream({
1594 write(value, c) {
1595 strictEqual(value, 1);
1596 c.error(new Error('boom'));
1597 },
1598 });
1599 
1600 const writer = ws.getWriter();
1601 
1602 // Should succeed
1603 await writer.write(1);
1604 
1605 await rejects(writer.closed, { message: 'boom' });
1606 }
1607 },
1608};
1609 
1610// Test new WritableStream() abort algorithm error handled
1611export const newWritableStreamAbortError = {
1612 async test() {
1613 // Sync abort error
1614 {
1615 const ws = new WritableStream({
1616 abort(reason) {
1617 strictEqual(reason, 1);
1618 throw new Error('boom');
1619 },
1620 });
1621 
1622 const writer = ws.getWriter();
1623 
1624 await rejects(writer.abort(1), { message: 'boom' });
1625 
1626 // Calling abort again returns the same rejected promise
1627 await rejects(writer.abort(1), { message: 'boom' });
1628 }
1629 
1630 // Async abort error
1631 {
1632 const ws = new WritableStream({
1633 async abort(reason) {
1634 strictEqual(reason, 1);
1635 throw new Error('boom');
1636 },
1637 });
1638 
1639 const writer = ws.getWriter();
1640 await rejects(writer.abort(1), { message: 'boom' });
1641 }
1642 
1643 // Abort with controller.error
1644 {
1645 let controller;
1646 const ws = new WritableStream({
1647 start(c) {
1648 controller = c;
1649 },
1650 async abort(reason) {
1651 strictEqual(reason, 1);
1652 controller.error(new Error('ignored'));
1653 },
1654 });
1655 
1656 const writer = ws.getWriter();
1657 
1658 await writer.abort(1);
1659 
1660 // The closed promise will use the abort reason, not the error
1661 // reported in the controller
1662 const results = await Promise.allSettled([writer.closed]);
1663 strictEqual(results[0].status, 'rejected');
1664 strictEqual(results[0].reason, 1);
1665 }
1666 },
1667};
1668 
1669// Test new WritableStream() close algorithm error handled
1670export const newWritableStreamCloseError = {
1671 async test() {
1672 // Sync close error
1673 {
1674 const ws = new WritableStream({
1675 close() {
1676 throw new Error('boom');
1677 },
1678 });
1679 
1680 const writer = ws.getWriter();
1681 await rejects(writer.close(), { message: 'boom' });
1682 }
1683 
1684 // Async close error
1685 {
1686 const ws = new WritableStream({
1687 async close() {
1688 throw new Error('boom');
1689 },
1690 });
1691 
1692 const writer = ws.getWriter();
1693 await rejects(writer.close(), { message: 'boom' });
1694 }
1695 
1696 // Close with controller.error (ignored)
1697 {
1698 let controller;
1699 const ws = new WritableStream({
1700 start(c) {
1701 controller = c;
1702 },
1703 async close() {
1704 controller.error(new Error('ignored'));
1705 },
1706 });
1707 
1708 const writer = ws.getWriter();
1709 
1710 // In this case, the error reported in the close algorithm is ignored
1711 await Promise.all([writer.close(), writer.closed]);
1712 }
1713 },
1714};
1715 
1716// Test WritableStream multiple pending writes allowed
1717export const writableStreamMultiplePendingWrites = {
1718 async test() {
1719 const expectedWrites = ['hello', 'there'];
1720 const ws = new WritableStream({
1721 async write(value) {
1722 await scheduler.wait(10);
1723 strictEqual(value, expectedWrites.shift());
1724 },
1725 });
1726 
1727 const writer = ws.getWriter();
1728 
1729 await Promise.all([writer.write('hello'), writer.write('there')]);
1730 },
1731};
1732 
1733// Test WritableStream writing a subarray works
1734export const writableStreamWriteSubarray = {
1735 async test() {
1736 const u8 = new Uint8Array([1, 2, 3, 4]);
1737 const sub = u8.subarray(1, 3);
1738 
1739 const ws = new WritableStream({
1740 write(value) {
1741 strictEqual(value, sub);
1742 },
1743 });
1744 
1745 const writer = ws.getWriter();
1746 
1747 await writer.write(sub);
1748 },
1749};
1750 
1751// Test WritableStream writing any javascript value works
1752export const writableStreamWriteAny = {
1753 async test() {
1754 // Make a copy since we'll shift() from it
1755 const expectedWrites = [
1756 'hello',
1757 true,
1758 1,
1759 1.1,
1760 undefined,
1761 NaN,
1762 Infinity,
1763 new Uint8Array(1),
1764 {},
1765 [],
1766 ];
1767 // Keep original for writing
1768 const valuesToWrite = [...expectedWrites];
1769 
1770 const ws = new WritableStream({
1771 async write(value) {
1772 await scheduler.wait(1);
1773 const expected = expectedWrites.shift();
1774 // Use Object.is for NaN and -0/+0 handling (same as testharness same_value)
1775 ok(Object.is(value, expected), `expected ${expected} but got ${value}`);
1776 },
1777 });
1778 
1779 const writer = ws.getWriter();
1780 
1781 await Promise.all(valuesToWrite.map((i) => writer.write(i)));
1782 },
1783};
1784 
1785// Test WritableStream desiredSize calculated correctly
1786export const writableStreamDesiredSize = {
1787 async test() {
1788 const ws = new WritableStream(
1789 {
1790 async write() {
1791 await scheduler.wait(10);
1792 },
1793 },
1794 {
1795 highWaterMark: 2,
1796 }
1797 );
1798 
1799 const writer = ws.getWriter();
1800 
1801 const firstReady = writer.ready;
1802 await firstReady;
1803 
1804 strictEqual(writer.desiredSize, 2);
1805 const write1 = writer.write(1);
1806 strictEqual(writer.desiredSize, 1);
1807 
1808 const write2 = writer.write(2);
1809 const write3 = writer.write(3);
1810 
1811 strictEqual(writer.desiredSize, -1);
1812 
1813 await Promise.all([write1, write2, write3]);
1814 
1815 ok(firstReady != writer.ready);
1816 await writer.ready;
1817 },
1818};
1819 
1820// Test WritableStream writes can be aborted
1821export const writableStreamWriteAbort = {
1822 async test() {
1823 let aborted = false;
1824 const ws = new WritableStream({
1825 start(c) {
1826 c.signal.addEventListener(
1827 'abort',
1828 () => {
1829 strictEqual(c.signal.reason.message, 'boom');
1830 aborted = true;
1831 },
1832 { once: true }
1833 );
1834 },
1835 async write() {
1836 await scheduler.wait(10);
1837 },
1838 });
1839 
1840 const writer = ws.getWriter();
1841 const write1 = writer.write(1);
1842 const write2 = writer.write(2);
1843 
1844 writer.abort(new Error('boom'));
1845 
1846 const [result1, result2] = await Promise.allSettled([write1, write2]);
1847 
1848 // Both writes fail
1849 strictEqual(result1.status, 'rejected');
1850 strictEqual(result2.status, 'rejected');
1851 strictEqual(result1.reason.message, 'boom');
1852 strictEqual(result2.reason.message, 'boom');
1853 
1854 ok(aborted);
1855 
1856 // Aborting puts the stream into a persistent errored state
1857 writer.releaseLock();
1858 const writer2 = ws.getWriter();
1859 
1860 await rejects(writer2.write('should not be allowed'), { message: 'boom' });
1861 },
1862};
1863 
1864// Test WritableStream uses size algorithm correctly
1865export const writableStreamSizeAlgorithm = {
1866 async test() {
1867 // Size algorithm called
1868 {
1869 let sizeCalled = false;
1870 const ws = new WritableStream(
1871 {
1872 async write() {
1873 await scheduler.wait(10);
1874 },
1875 },
1876 {
1877 highWaterMark: 2,
1878 size(value) {
1879 sizeCalled = true;
1880 strictEqual(value, 'hello');
1881 return 2;
1882 },
1883 }
1884 );
1885 
1886 const writer = ws.getWriter();
1887 strictEqual(writer.desiredSize, 2);
1888 const write = writer.write('hello');
1889 ok(sizeCalled);
1890 strictEqual(writer.desiredSize, 0);
1891 await write;
1892 }
1893 
1894 // Size algorithm error
1895 {
1896 const ws = new WritableStream(
1897 {},
1898 {
1899 size() {
1900 throw new Error('boom');
1901 },
1902 }
1903 );
1904 
1905 const writer = ws.getWriter();
1906 await rejects(writer.write('hello'), { message: 'boom' });
1907 }
1908 },
1909};
1910 
1911// Test WritableStream aborting should replace ready promise
1912export const writableStreamAbortReadyRejected = {
1913 async test() {
1914 const ws = new WritableStream();
1915 const writer = ws.getWriter();
1916 const ready = writer.ready;
1917 
1918 await ready;
1919 
1920 writer.abort('boom');
1921 
1922 ok(ready !== writer.ready);
1923 
1924 const results = await Promise.allSettled([writer.ready]);
1925 strictEqual(results[0].status, 'rejected');
1926 strictEqual(results[0].reason, 'boom');
1927 },
1928};
1929 
1930// Test WritableStream abort with no argument defaults to undefined
1931export const writableStreamAbortOptional = {
1932 async test() {
1933 const ws = new WritableStream();
1934 const writer = ws.getWriter();
1935 await writer.abort();
1936 
1937 const results = await Promise.allSettled([writer.closed]);
1938 strictEqual(results[0].status, 'rejected');
1939 strictEqual(results[0].reason, undefined);
1940 },
1941};
1942 
1943// Test WritableStream abort while starting rejects ready promise
1944export const writableStreamAbortWhileStarting = {
1945 async test() {
1946 const ws = new WritableStream({
1947 async start() {},
1948 });
1949 
1950 const writer = ws.getWriter();
1951 writer.abort('boom');
1952 
1953 const results = await Promise.allSettled([writer.ready]);
1954 strictEqual(results[0].status, 'rejected');
1955 strictEqual(results[0].reason, 'boom');
1956 },
1957};
1958 
1959// Test WritableStream abort error while writing rejects abort promise
1960export const writableStreamAbortWhileWriting = {
1961 async test() {
1962 const ws = new WritableStream({
1963 async write() {
1964 await scheduler.wait(10);
1965 },
1966 
1967 async abort() {
1968 throw new Error('boom');
1969 },
1970 });
1971 
1972 const writer = ws.getWriter();
1973 const write = writer.write('test');
1974 
1975 await rejects(writer.abort(), { message: 'boom' });
1976 
1977 await write;
1978 
1979 // The write should reject with undefined because the writer.abort() above
1980 // specified undefined. The abort algorithm should not be called again.
1981 await rejects(writer.write('should fail'), (err) => err === undefined);
1982 },
1983};
1984 
1985// Test WritableStream releaseLock while aborting should reject closed promise
1986export const writableStreamReleaseLockWhileAborting = {
1987 async test() {
1988 const ws = new WritableStream({
1989 async write() {
1990 await scheduler.wait(10);
1991 },
1992 });
1993 
1994 const writer = ws.getWriter();
1995 writer.write('test');
1996 
1997 writer.abort();
1998 const closed = writer.closed;
1999 
2000 writer.releaseLock();
2001 
2002 await rejects(closed, TypeError);
2003 },
2004};
2005 
2006// Test WritableStream throw during in flight close rejects abort and closed promise
2007export const writableStreamCloseThrowRejectsPromises = {
2008 async test() {
2009 const ws = new WritableStream({
2010 async close() {
2011 throw new Error('boom');
2012 },
2013 });
2014 
2015 const writer = ws.getWriter();
2016 const close = writer.close();
2017 const abort = writer.abort();
2018 const closed = writer.closed;
2019 
2020 const res = await Promise.allSettled([close, abort, closed]);
2021 
2022 strictEqual(res[0].status, 'rejected');
2023 strictEqual(res[1].status, 'rejected');
2024 strictEqual(res[2].status, 'rejected');
2025 
2026 // The close and abort promises are rejected with the error thrown,
2027 // the closed promise, however, should be rejected with the reason
2028 // given in the abort(), which is defaulted to undefined.
2029 strictEqual(res[0].reason.message, 'boom');
2030 strictEqual(res[1].reason.message, 'boom');
2031 strictEqual(res[2].reason, undefined);
2032 },
2033};
2034 
2035// Test WritableStream sink abort not called while write or close in flight
2036export const writableStreamAbortTiming = {
2037 async test() {
2038 // Abort waits for start
2039 {
2040 let started = false;
2041 const ws = new WritableStream({
2042 async start() {
2043 await scheduler.wait(10);
2044 started = true;
2045 },
2046 abort() {
2047 ok(started, 'The stream should have started first');
2048 },
2049 });
2050 
2051 await ws.abort();
2052 }
2053 
2054 // Abort waits for write
2055 {
2056 let writeCompleted = false;
2057 const ws = new WritableStream({
2058 async write() {
2059 await scheduler.wait(10);
2060 writeCompleted = true;
2061 },
2062 abort() {
2063 ok(writeCompleted, 'The write should have completed');
2064 },
2065 });
2066 
2067 const writer = ws.getWriter();
2068 const write = writer.write('hello');
2069 const abort = writer.abort();
2070 
2071 await Promise.allSettled([write, abort]);
2072 }
2073 
2074 // Abort waits for close
2075 {
2076 let closeCompleted = false;
2077 const ws = new WritableStream({
2078 async close() {
2079 await scheduler.wait(10);
2080 closeCompleted = true;
2081 },
2082 abort() {
2083 ok(closeCompleted, 'The close should have completed');
2084 },
2085 });
2086 
2087 const writer = ws.getWriter();
2088 const close = writer.close();
2089 const abort = writer.abort();
2090 
2091 await Promise.allSettled([close, abort]);
2092 }
2093 },
2094};
2095 
2096// Test WritableStream abort during write should trigger abort algorithm with close pending
2097export const writableStreamAbortWriteClosePending = {
2098 async test() {
2099 let abortCalled = false;
2100 const ws = new WritableStream({
2101 async write() {
2102 await scheduler.wait(10);
2103 },
2104 abort() {
2105 abortCalled = true;
2106 },
2107 });
2108 
2109 const writer = ws.getWriter();
2110 const write = writer.write('hello');
2111 const close = writer.close();
2112 const abort = writer.abort();
2113 
2114 const res = await Promise.allSettled([write, close, abort]);
2115 ok(abortCalled);
2116 
2117 strictEqual(res[0].status, 'fulfilled'); // Write finishes
2118 strictEqual(res[1].status, 'rejected'); // Pending close is aborted
2119 strictEqual(res[2].status, 'fulfilled'); // Abort finishes
2120 },
2121};
2122 
2123// Test WritableStream ready promise rejects on controller error not waiting for in flight write
2124export const writableStreamErrorDuringInFlightWrite = {
2125 async test() {
2126 let controller;
2127 const ws = new WritableStream({
2128 start(c) {
2129 controller = c;
2130 },
2131 async write() {
2132 await scheduler.wait(10);
2133 },
2134 });
2135 
2136 const writer = ws.getWriter();
2137 
2138 const write = writer.write('hello').catch(() => {});
2139 
2140 controller.error('boom');
2141 
2142 await Promise.all([write, writer.ready.catch(() => {})]);
2143 },
2144};
2145 
2146// Test WritableStream start errors after abort, close rejects
2147export const writableStreamStartErrorAfterAbort = {
2148 async test() {
2149 const ws = new WritableStream({
2150 async start() {
2151 await scheduler.wait(10);
2152 throw new Error('boom');
2153 },
2154 });
2155 
2156 ws.abort();
2157 
2158 // close() should return a rejected promise. The abort was called
2159 // before start finished, so the stream enters an errored state.
2160 await ws.close().catch(() => {});
2161 
2162 // Verify the stream is errored
2163 strictEqual(ws.locked, false);
2164 },
2165};
2166 
2167// Test WritableStream rejected sink write does not prevent sink abort
2168export const writableStreamRejectedWriteNoPreventAbort = {
2169 async test() {
2170 let abortCalled = false;
2171 const ws = new WritableStream({
2172 write() {
2173 throw new Error('boom');
2174 },
2175 abort() {
2176 abortCalled = true;
2177 },
2178 });
2179 
2180 const writer = ws.getWriter();
2181 const write = writer.write('hello');
2182 const abort = writer.abort();
2183 
2184 const res = await Promise.allSettled([write, abort]);
2185 ok(abortCalled);
2186 strictEqual(res[0].status, 'rejected');
2187 strictEqual(res[1].status, 'fulfilled');
2188 },
2189};
2190 
2191// Test WritableStream aborted twice
2192export const writableStreamAbortedTwice = {
2193 async test() {
2194 const ws = new WritableStream();
2195 const abort1 = ws.abort();
2196 const abort2 = ws.abort();
2197 
2198 await Promise.all([abort1, abort2]);
2199 },
2200};
2201 
2202// Test WritableStream aborting errored stream rejects with stored error
2203export const writableStreamAbortOnErroredResolves = {
2204 async test() {
2205 const ws = new WritableStream({
2206 start(c) {
2207 c.error(new Error('boom'));
2208 },
2209 });
2210 
2211 // When aborting an already-errored stream, abort() rejects with the stored error
2212 await rejects(ws.abort(), Error);
2213 },
2214};
2215 
2216// Test WritableStream sink abort not called if stream errored before abort
2217export const writableStreamSinkAlgNoCallErrorBeforeAbort = {
2218 test() {
2219 let controller;
2220 let abortCalled = false;
2221 const ws = new WritableStream({
2222 start(c) {
2223 controller = c;
2224 },
2225 abort() {
2226 abortCalled = true;
2227 },
2228 });
2229 
2230 controller.error(new Error('boom'));
2231 ws.abort(new Error('bang'));
2232 
2233 ok(!abortCalled);
2234 },
2235};
2236 
2237// Test WritableStream writer with pending abort, ready should reject
2238export const writableStreamWriterWithPendingAbort = {
2239 async test() {
2240 const ws = new WritableStream();
2241 ws.abort(new Error('boom'));
2242 const writer = ws.getWriter();
2243 
2244 await rejects(writer.ready, Error);
2245 },
2246};
2247 
2248// Test WritableStream promises resolved in order
2249export const writableStreamPromisesResolvedInOrder = {
2250 async test() {
2251 // Write before close
2252 {
2253 let closeFinished = false;
2254 const ws = new WritableStream();
2255 
2256 const writer = ws.getWriter();
2257 
2258 const write = writer.write('hello').then(() => ok(!closeFinished));
2259 const close = writer.close();
2260 const closed = writer.closed.then(() => (closeFinished = true));
2261 
2262 // Closed promise should not resolve before fulfilled write.
2263 await Promise.allSettled([write, close, closed]);
2264 }
2265 
2266 // Rejected write before close
2267 {
2268 let closeFinished = false;
2269 let writeFailed = false;
2270 const ws = new WritableStream({
2271 write() {
2272 throw new Error('boom');
2273 },
2274 });
2275 
2276 const writer = ws.getWriter();
2277 
2278 const write = writer.write('hello').catch(() => {
2279 ok(!closeFinished);
2280 writeFailed = true;
2281 });
2282 const close = writer.close();
2283 const closed = writer.closed.then(() => (closeFinished = true));
2284 
2285 // Closed promise should not resolve before rejected write.
2286 await Promise.allSettled([write, close, closed]);
2287 ok(writeFailed);
2288 }
2289 
2290 // Writes resolved in order when aborting
2291 {
2292 const order = [];
2293 const ws = new WritableStream({
2294 async write() {
2295 await scheduler.wait(10);
2296 },
2297 });
2298 
2299 const writer = ws.getWriter();
2300 
2301 const write1 = writer.write('hello').then(() => order.push(1));
2302 const write2 = writer.write('hello').catch(() => order.push(2));
2303 const write3 = writer.write('hello').catch(() => order.push(3));
2304 const abort = writer.abort();
2305 
2306 await Promise.allSettled([write1, write2, write3, abort]);
2307 
2308 strictEqual(order[0], 1);
2309 strictEqual(order[1], 2);
2310 strictEqual(order[2], 3);
2311 }
2312 },
2313};
2314 
2315// =====================================================================================
2316// Misc tests
2317// =====================================================================================
2318 
2319// Test highWaterMark validation
2320export const highWaterMarkValidated = {
2321 test() {
2322 [-1, -Infinity, NaN, {}, 'foo'].forEach((highWaterMark) => {
2323 throws(() => new WritableStream(undefined, { highWaterMark }), TypeError);
2324 throws(() => new ReadableStream(undefined, { highWaterMark }), TypeError);
2325 });
2326 },
2327};
2328 
2329// Test QueuingStrategy objects work
2330export const queuingStrategies = {
2331 test() {
2332 // ByteLengthQueuingStrategy
2333 {
2334 const strategy = new ByteLengthQueuingStrategy({ highWaterMark: 10 });
2335 
2336 let startRan = false;
2337 
2338 // Make sure we can create a stream using the strategy without error.
2339 new ReadableStream(
2340 {
2341 start(c) {
2342 strictEqual(c.desiredSize, 10);
2343 c.enqueue(new Uint8Array(2));
2344 strictEqual(c.desiredSize, 8);
2345 startRan = true;
2346 },
2347 },
2348 strategy
2349 );
2350 
2351 const { highWaterMark, size } = strategy;
2352 
2353 ok(startRan);
2354 strictEqual(highWaterMark, 10);
2355 strictEqual(size('nothing'), undefined);
2356 strictEqual(size(123), undefined);
2357 strictEqual(size(undefined), undefined);
2358 strictEqual(size(null), undefined);
2359 strictEqual(size(), undefined);
2360 strictEqual(size(new ArrayBuffer(10)), 10);
2361 strictEqual(size(new Uint8Array(10)), 10);
2362 }
2363 
2364 // CountQueuingStrategy
2365 {
2366 const strategy = new CountQueuingStrategy({ highWaterMark: 9 });
2367 
2368 let startRan = false;
2369 
2370 // Make sure we can create a stream using the strategy without error.
2371 new ReadableStream(
2372 {
2373 start(c) {
2374 strictEqual(c.desiredSize, 9);
2375 c.enqueue(new Uint8Array(2));
2376 strictEqual(c.desiredSize, 8);
2377 startRan = true;
2378 },
2379 },
2380 strategy
2381 );
2382 
2383 const { highWaterMark, size } = strategy;
2384 
2385 ok(startRan);
2386 strictEqual(highWaterMark, 9);
2387 strictEqual(size('nothing'), 1);
2388 strictEqual(size(123), 1);
2389 strictEqual(size(undefined), 1);
2390 strictEqual(size(null), 1);
2391 strictEqual(size(), 1);
2392 strictEqual(size(new ArrayBuffer(10)), 1);
2393 strictEqual(size(new Uint8Array(10)), 1);
2394 }
2395 },
2396};
2397 
2398// Test proper default highwater mark
2399export const hwmDefault = {
2400 async test() {
2401 let pulled = 0;
2402 new ReadableStream({
2403 start(c) {
2404 strictEqual(c.desiredSize, 1);
2405 },
2406 pull() {
2407 pulled++;
2408 },
2409 });
2410 
2411 new ReadableStream({
2412 type: 'bytes',
2413 start(c) {
2414 strictEqual(c.desiredSize, 0);
2415 },
2416 pull() {
2417 pulled += 2;
2418 },
2419 });
2420 
2421 await scheduler.wait(1);
2422 strictEqual(pulled, 1);
2423 },
2424};
2425 
2426// Test byobreader regression
2427export const byobreaderRegression = {
2428 async test() {
2429 function newReadableStream(chunks) {
2430 chunks = chunks.filter((ch) => ch !== 0);
2431 return new ReadableStream({
2432 type: 'bytes',
2433 start(c) {
2434 if (chunks.length === 0) {
2435 c.close();
2436 }
2437 },
2438 pull(c) {
2439 c.enqueue(new Uint8Array(chunks.shift()));
2440 if (chunks.length === 0) {
2441 c.close();
2442 }
2443 },
2444 });
2445 }
2446 
2447 const rs = newReadableStream([]);
2448 // Ensure that getting a byob reader on a closed byte stream works correctly.
2449 const reader = rs.getReader({ mode: 'byob' });
2450 const { done } = await reader.read(new Uint8Array(10));
2451 ok(done);
2452 },
2453};
2454 
2455// Test writer double close
2456export const writerDoubleClose = {
2457 async test() {
2458 const ws = new WritableStream({
2459 write() {},
2460 });
2461 const writer = ws.getWriter();
2462 
2463 writer.write(123);
2464 
2465 writer.close();
2466 // With capture_async_api_throws, async methods return rejected promises instead of throwing
2467 if (Cloudflare.compatibilityFlags.capture_async_api_throws) {
2468 await rejects(writer.close(), TypeError);
2469 } else {
2470 throws(() => writer.close(), TypeError);
2471 }
2472 },
2473};
2474 
2475// =====================================================================================
2476// GC tests (these require --expose-gc v8 flag)
2477// =====================================================================================
2478 
2479// Test ReadableStream object references are held through gc
2480export const readableStreamReferencesHold = {
2481 async test() {
2482 let controller;
2483 let reader;
2484 let read;
2485 
2486 // Byte stream
2487 {
2488 const rs = new ReadableStream({
2489 type: 'bytes',
2490 start(c) {
2491 controller = c;
2492 },
2493 });
2494 
2495 reader = rs.getReader({ mode: 'byob' });
2496 }
2497 
2498 await scheduler.wait(10);
2499 gc();
2500 
2501 {
2502 read = reader.read(new Uint8Array(1));
2503 reader = undefined;
2504 }
2505 
2506 await scheduler.wait(10);
2507 gc();
2508 
2509 {
2510 controller.enqueue(new Uint8Array([1]));
2511 controller = undefined;
2512 const { value, done } = await read;
2513 ok(!done);
2514 strictEqual(value[0], 1);
2515 }
2516 
2517 // Value stream
2518 {
2519 let controller;
2520 let reader;
2521 let read;
2522 
2523 {
2524 const rs = new ReadableStream({
2525 start(c) {
2526 controller = c;
2527 },
2528 });
2529 reader = rs.getReader();
2530 }
2531 
2532 await scheduler.wait(10);
2533 gc();
2534 
2535 {
2536 read = reader.read();
2537 reader = undefined;
2538 }
2539 
2540 await scheduler.wait(10);
2541 gc();
2542 
2543 {
2544 controller.enqueue('hello');
2545 controller = undefined;
2546 const { value, done } = await read;
2547 ok(!done);
2548 strictEqual(value, 'hello');
2549 }
2550 }
2551 },
2552};
2553 
2554// Test WritableStream object references are held through gc
2555export const writableStreamGc = {
2556 async test() {
2557 let controller;
2558 let writer;
2559 let write;
2560 
2561 {
2562 const ws = new WritableStream({
2563 start(c) {
2564 controller = c;
2565 },
2566 });
2567 writer = ws.getWriter();
2568 }
2569 
2570 await scheduler.wait(10);
2571 gc();
2572 
2573 {
2574 write = writer.write(1);
2575 writer = undefined;
2576 }
2577 
2578 await scheduler.wait(10);
2579 gc();
2580 
2581 {
2582 await write;
2583 strictEqual(controller.signal.aborted, false);
2584 }
2585 },
2586};
2587 
2588// Test ReadableStream with async iterator gc works
2589export const asyncIteratorGc = {
2590 async test() {
2591 // This test verifies that the ReadableStream and its async iterator
2592 // are properly handled through gc
2593 function getNextPromise() {
2594 let values = new ReadableStream({
2595 async pull(controller) {
2596 await scheduler.wait(50);
2597 controller.enqueue('A');
2598 controller.close();
2599 },
2600 }).values();
2601 values.next();
2602 const promise = values.next();
2603 values = undefined;
2604 return promise;
2605 }
2606 
2607 let promise = getNextPromise();
2608 gc();
2609 strictEqual((await promise).done, true);
2610 promise = undefined;
2611 gc();
2612 },
2613};