Skip to content
File

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

javascript792 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 
5import { strictEqual, rejects } from 'node:assert';
6 
7// Test Response body methods with JS-backed BYOB ReadableStream
8export const responseBodyMethodsJsByob = {
9 async test() {
10 const enc = new TextEncoder();
11 const dec = new TextDecoder();
12 
13 {
14 const rs = new ReadableStream({
15 type: 'bytes',
16 async pull(c) {
17 if (c.byobRequest) {
18 enc.encodeInto('hello', c.byobRequest.view);
19 c.byobRequest.respond(5);
20 c.close();
21 }
22 },
23 });
24 
25 const resp = new Response(rs);
26 
27 strictEqual(dec.decode(await resp.arrayBuffer()), 'hello');
28 }
29 
30 {
31 const rs = new ReadableStream({
32 type: 'bytes',
33 async pull(c) {
34 if (c.byobRequest) {
35 enc.encodeInto('hello', c.byobRequest.view);
36 c.byobRequest.respond(5);
37 c.close();
38 }
39 },
40 });
41 
42 const resp = new Response(rs);
43 
44 strictEqual(await resp.text(), 'hello');
45 }
46 },
47};
48 
49// Test Request body methods with JS-backed ReadableStream
50export const requestBodyMethodsJsByob = {
51 async test() {
52 const enc = new TextEncoder();
53 const wrapped = new ReadableStream({
54 type: 'bytes',
55 async pull(c) {
56 c.enqueue(enc.encode('hello'));
57 c.close();
58 },
59 });
60 
61 const req = new Request('http://example.com', {
62 method: 'POST',
63 body: wrapped,
64 });
65 const text = await req.text();
66 
67 strictEqual(text, 'hello');
68 },
69};
70 
71// Test basic JS ReadableStream as Response body
72export const jsSource = {
73 async test() {
74 const enc = new TextEncoder();
75 const rs = new ReadableStream({
76 start(c) {
77 c.enqueue(enc.encode('hello'));
78 c.close();
79 },
80 });
81 
82 const response = new Response(rs);
83 strictEqual(await response.text(), 'hello');
84 },
85};
86 
87// Test JS ReadableStream with async pull as Response body
88export const jsSourceAsyncPull = {
89 async test() {
90 const enc = new TextEncoder();
91 const chunks = [enc.encode('hello'), enc.encode('there')];
92 const rs = new ReadableStream({
93 async pull(c) {
94 await scheduler.wait(10);
95 if (chunks.length > 0) c.enqueue(chunks.shift());
96 if (chunks.length === 0) c.close();
97 },
98 });
99 
100 const response = new Response(rs);
101 strictEqual(await response.text(), 'hellothere');
102 },
103};
104 
105// Test BYOB ReadableStream as Response body
106export const jsByteSource = {
107 async test() {
108 const enc = new TextEncoder();
109 const rs = new ReadableStream({
110 type: 'bytes',
111 pull(c) {
112 const request = c.byobRequest;
113 if (request != null) {
114 enc.encodeInto('hello', request.view);
115 request.respond(5);
116 c.close();
117 }
118 },
119 });
120 
121 const response = new Response(rs);
122 strictEqual(await response.text(), 'hello');
123 },
124};
125 
126// Test BYOB ReadableStream with multiple chunks
127export const jsByteSourceMultipleChunks = {
128 async test() {
129 const enc = new TextEncoder();
130 const chunks = ['hello', 'there', 'this', 'is', 'a', 'test'];
131 
132 const rs = new ReadableStream({
133 type: 'bytes',
134 async pull(c) {
135 await scheduler.wait(10);
136 const request = c.byobRequest;
137 if (request != null) {
138 const chunk = chunks.shift();
139 if (chunk !== undefined) {
140 const { written } = enc.encodeInto(chunk, request.view);
141 request.respond(written);
142 } else {
143 c.close();
144 request.respond(0);
145 }
146 }
147 },
148 });
149 
150 const response = new Response(rs);
151 strictEqual(await response.text(), 'hellotherethisisatest');
152 },
153};
154 
155// Test teed JS ReadableStream as Response body
156export const jsTeeSource = {
157 async test() {
158 const enc = new TextEncoder();
159 const dec = new TextDecoder();
160 const readable = new ReadableStream({
161 start(c) {
162 c.enqueue(enc.encode('hello'));
163 c.close();
164 },
165 });
166 
167 const [branch1, branch2] = readable.tee();
168 const reader = branch2.getReader();
169 
170 // By reading first, we ensure that the branch1 will have queued
171 // read data before the Response is actually transmitted. Both tee
172 // branches should still resolve to the same data.
173 
174 const result = await reader.read();
175 strictEqual(dec.decode(result.value), 'hello');
176 
177 const response = new Response(branch1);
178 strictEqual(await response.text(), 'hello');
179 },
180};
181 
182// Test closed teed BYOB stream
183export const jsTeeClose = {
184 async test() {
185 const rs = new ReadableStream({
186 type: 'bytes',
187 async pull(c) {
188 if (c.byobRequest) {
189 c.close();
190 c.byobRequest.respond(0);
191 }
192 },
193 });
194 
195 const [branch] = rs.tee();
196 const response = new Response(branch);
197 strictEqual(await response.text(), '');
198 },
199};
200 
201// Test large enqueue with synchronous close (default stream)
202export const bigEnqueue = {
203 async test() {
204 const a = 'a'.repeat(4096 * 2);
205 const enc = new TextEncoder();
206 const rs = new ReadableStream({
207 start(c) {
208 c.enqueue(enc.encode(a));
209 c.close();
210 },
211 });
212 
213 const response = new Response(rs);
214 const text = await response.text();
215 strictEqual(text.length, 8192);
216 strictEqual(text, a);
217 },
218};
219 
220// Test large enqueue with synchronous close (bytes stream)
221export const bigEnqueueBytes = {
222 async test() {
223 const a = 'a'.repeat(4096 * 2);
224 const enc = new TextEncoder();
225 const rs = new ReadableStream({
226 type: 'bytes',
227 start(c) {
228 c.enqueue(enc.encode(a));
229 c.close();
230 },
231 });
232 
233 const response = new Response(rs);
234 const text = await response.text();
235 strictEqual(text.length, 8192);
236 strictEqual(text, a);
237 },
238};
239 
240// Test large enqueue via IdentityTransformStream (sync close)
241export const bigEnqueueViaIdentityTransform = {
242 async test() {
243 const a = 'a'.repeat(4096 * 2 + 345);
244 const enc = new TextEncoder();
245 
246 const rs = new ReadableStream({
247 start(c) {
248 c.enqueue(enc.encode(a));
249 c.close();
250 },
251 });
252 
253 const transform = new IdentityTransformStream();
254 const response = new Response(rs.pipeThrough(transform));
255 
256 const text = await response.text();
257 strictEqual(text.length, 8537);
258 strictEqual(text, a);
259 },
260};
261 
262// Test large enqueue via IdentityTransformStream (async close)
263export const bigEnqueueViaIdentityTransformAsync = {
264 async test() {
265 const a = 'a'.repeat(4096 * 2 + 345);
266 const enc = new TextEncoder();
267 
268 const rs = new ReadableStream({
269 async start(c) {
270 c.enqueue(enc.encode(a));
271 await scheduler.wait(10);
272 c.close();
273 },
274 });
275 
276 const transform = new IdentityTransformStream();
277 const response = new Response(rs.pipeThrough(transform));
278 
279 const text = await response.text();
280 strictEqual(text.length, 8537);
281 strictEqual(text, a);
282 },
283};
284 
285// Test large enqueue via JS TransformStream (sync close)
286export const bigEnqueueViaJsTransform = {
287 async test() {
288 const a = 'a'.repeat(4096 * 2 + 345);
289 const enc = new TextEncoder();
290 
291 const rs = new ReadableStream({
292 start(c) {
293 c.enqueue(enc.encode(a));
294 c.close();
295 },
296 });
297 
298 const transform = new TransformStream();
299 const response = new Response(rs.pipeThrough(transform));
300 
301 const text = await response.text();
302 strictEqual(text.length, 8537);
303 strictEqual(text, a);
304 },
305};
306 
307// Test large enqueue via JS TransformStream (async close)
308export const bigEnqueueViaJsTransformAsync = {
309 async test() {
310 const a = 'a'.repeat(4096 * 2 + 345);
311 const enc = new TextEncoder();
312 
313 const rs = new ReadableStream({
314 async start(c) {
315 c.enqueue(enc.encode(a));
316 await scheduler.wait(10);
317 c.close();
318 },
319 });
320 
321 const transform = new TransformStream();
322 const response = new Response(rs.pipeThrough(transform));
323 
324 const text = await response.text();
325 strictEqual(text.length, 8537);
326 strictEqual(text, a);
327 },
328};
329 
330export const bigEnqueueViaJsTransformSplit = {
331 async test() {
332 const enc = new TextEncoder();
333 const dec = new TextDecoder();
334 
335 const rs = new ReadableStream({
336 start(c) {
337 c.enqueue(enc.encode('a'.repeat(4090)));
338 },
339 pull(c) {
340 c.enqueue(enc.encode('a'.repeat(4096 + 345)));
341 c.close();
342 },
343 });
344 
345 const transform = new TransformStream({
346 transform(chunk, controller) {
347 controller.enqueue(enc.encode(dec.decode(chunk).toUpperCase()));
348 },
349 });
350 
351 const response = new Response(rs.pipeThrough(transform));
352 
353 const text = await response.text();
354 strictEqual(text.length, 4090 + 4096 + 345);
355 strictEqual(text, 'A'.repeat(4090 + 4096 + 345));
356 },
357};
358 
359// Test enqueue same chunk multiple times (default stream)
360export const enqueueChunkMultipleTimes = {
361 async test() {
362 const chunk = new TextEncoder().encode('ping!');
363 const rs = new ReadableStream({
364 start(controller) {
365 controller.enqueue(chunk);
366 controller.enqueue(chunk);
367 controller.enqueue(chunk);
368 controller.enqueue(chunk);
369 controller.close();
370 },
371 });
372 
373 const response = new Response(rs);
374 strictEqual(await response.text(), 'ping!ping!ping!ping!');
375 },
376};
377 
378// Test enqueue same chunk multiple times errors in bytes stream
379export const enqueueChunkMultipleTimesBytes = {
380 async test() {
381 const chunk = new TextEncoder().encode('ping!');
382 const rs = new ReadableStream({
383 type: 'bytes',
384 start(controller) {
385 controller.enqueue(chunk);
386 try {
387 controller.enqueue(chunk);
388 throw new Error('this should have failed because chunk is size 0');
389 } catch (err) {
390 if (err.message !== 'Cannot enqueue a zero-length ArrayBuffer.') {
391 throw new Error('Incorrect error: ' + err.message);
392 }
393 controller.close();
394 }
395 },
396 });
397 
398 const response = new Response(rs);
399 strictEqual(await response.text(), 'ping!');
400 },
401};
402 
403export const bigEnqueueOddSize = {
404 async test() {
405 const a = 'a'.repeat(4096 * 2 + 345);
406 const enc = new TextEncoder();
407 const rs = new ReadableStream({
408 start(c) {
409 c.enqueue(enc.encode(a));
410 c.close();
411 },
412 });
413 
414 const response = new Response(rs);
415 const text = await response.text();
416 strictEqual(text.length, 8537);
417 strictEqual(text, a);
418 },
419};
420 
421// In this test, we write data into an IdentityTransformStream
422// We then read from that in a JS ReadableStream
423// We then pipe that through a JS TransformStream
424// We then respond with the TransformStream's readable.
425// We use parallel writes to simulate waitUntil() behavior.
426export const multistepTransform = {
427 async test() {
428 const enc = new TextEncoder();
429 const dec = new TextDecoder();
430 
431 const { readable, writable } = new IdentityTransformStream();
432 const reader = readable.getReader({ mode: 'byob' });
433 
434 const rs = new ReadableStream({
435 type: 'bytes',
436 start(c) {
437 c.enqueue(enc.encode('bbbb'));
438 },
439 async pull(c) {
440 const buffer = new Uint8Array(4096);
441 const result = await reader.read(buffer);
442 if (result.done) {
443 c.enqueue(enc.encode('bye'));
444 c.close();
445 } else {
446 const view = c.byobRequest.view;
447 const toCopy = Math.min(result.value.byteLength, view.byteLength);
448 new Uint8Array(view.buffer, view.byteOffset, toCopy).set(
449 result.value.subarray(0, toCopy)
450 );
451 c.byobRequest.respond(toCopy);
452 }
453 },
454 });
455 
456 const transform = new TransformStream({
457 transform(chunk, controller) {
458 controller.enqueue(enc.encode(dec.decode(chunk).toUpperCase()));
459 },
460 });
461 
462 const response = new Response(rs.pipeThrough(transform));
463 
464 const writer = writable.getWriter();
465 const writePromise = (async () => {
466 await writer.write(enc.encode('a'.repeat(4090)));
467 await writer.write(enc.encode('a'.repeat(4096)));
468 await writer.write(enc.encode('a'.repeat(345)));
469 await writer.close();
470 })();
471 
472 const [text] = await Promise.all([response.text(), writePromise]);
473 
474 strictEqual(text.startsWith('BBBB'), true);
475 strictEqual(text.endsWith('BYE'), true);
476 strictEqual(text.length, 4 + 4090 + 4096 + 345 + 3);
477 },
478};
479 
480// In this test, we write data into an IdentityTransformStream
481// We then read from that in a JS ReadableStream
482// We then pipe that through a JS TransformStream with preventClose = true
483// When the pipe is done, we write a final chunk to the JS TransformStream
484// We then respond with the TransformStream's readable.
485export const multistepTransformPreventClose = {
486 async test() {
487 const enc = new TextEncoder();
488 const dec = new TextDecoder();
489 
490 const { readable, writable } = new IdentityTransformStream();
491 const reader = readable.getReader({ mode: 'byob' });
492 
493 const detachesBuffer =
494 Cloudflare.compatibilityFlags.streams_byob_reader_detaches_buffer;
495 
496 const rs = new ReadableStream({
497 type: 'bytes',
498 start(c) {
499 c.enqueue(enc.encode('bbbb'));
500 },
501 async pull(c) {
502 if (detachesBuffer) {
503 const buffer = new Uint8Array(4096);
504 const result = await reader.read(buffer);
505 if (result.done) {
506 c.enqueue(enc.encode('bye'));
507 c.close();
508 } else {
509 const view = c.byobRequest.view;
510 const toCopy = Math.min(result.value.byteLength, view.byteLength);
511 new Uint8Array(view.buffer, view.byteOffset, toCopy).set(
512 result.value.subarray(0, toCopy)
513 );
514 c.byobRequest.respond(toCopy);
515 }
516 } else {
517 const result = await reader.read(c.byobRequest.view);
518 if (result.done) {
519 c.enqueue(enc.encode('bye'));
520 c.close();
521 } else {
522 c.byobRequest.respondWithNewView(result.value);
523 }
524 }
525 },
526 });
527 
528 const transform = new TransformStream({
529 transform(chunk, controller) {
530 controller.enqueue(enc.encode(dec.decode(chunk).toUpperCase()));
531 },
532 });
533 
534 const promise = rs.pipeTo(transform.writable, { preventClose: true });
535 
536 const writer = writable.getWriter();
537 const writePromise = (async () => {
538 await writer.write(enc.encode('a'.repeat(4090)));
539 await writer.write(enc.encode('a'.repeat(4096)));
540 await writer.write(enc.encode('a'.repeat(345)));
541 await writer.close();
542 })();
543 
544 async function finishChunk() {
545 await promise;
546 const transformWriter = transform.writable.getWriter();
547 await transformWriter.write(enc.encode('all done'));
548 await transformWriter.close();
549 }
550 
551 const response = new Response(transform.readable);
552 
553 const [text] = await Promise.all([
554 response.text(),
555 writePromise,
556 finishChunk(),
557 ]);
558 
559 strictEqual(text.startsWith('BBBB'), true);
560 strictEqual(text.endsWith('ALL DONE'), true);
561 // 4 (BBBB) + 4090 + 4096 + 345 (A's) + 3 (BYE) + 8 (ALL DONE)
562 strictEqual(text.length, 4 + 4090 + 4096 + 345 + 3 + 8);
563 },
564};
565 
566export const jsSourceError = {
567 async test() {
568 const rs = new ReadableStream({
569 start(c) {
570 throw new Error('boom');
571 },
572 });
573 
574 const response = new Response(rs);
575 await rejects(response.text(), { name: 'Error', message: 'boom' });
576 },
577};
578 
579export const jsSourceErrorAsync = {
580 async test() {
581 const rs = new ReadableStream({
582 async pull(c) {
583 await scheduler.wait(10);
584 throw new Error('boom');
585 },
586 });
587 
588 const response = new Response(rs);
589 await rejects(response.text(), { name: 'Error', message: 'boom' });
590 },
591};
592 
593export const jsErroredSourceAsync = {
594 async test() {
595 const rs = new ReadableStream({
596 async start(c) {
597 throw new Error('boom');
598 },
599 });
600 
601 const response = new Response(rs);
602 await rejects(response.text(), { name: 'Error', message: 'boom' });
603 },
604};
605 
606export const jsErroredSourceAsyncDelayed = {
607 async test() {
608 const rs = new ReadableStream({
609 async start(c) {
610 await scheduler.wait(10);
611 throw new Error('boom');
612 },
613 });
614 
615 const response = new Response(rs);
616 await rejects(response.text(), { name: 'Error', message: 'boom' });
617 },
618};
619 
620export const jsNotBytesInPull = {
621 async test() {
622 const rs = new ReadableStream({
623 pull(c) {
624 c.enqueue('hello');
625 c.close();
626 },
627 });
628 
629 const response = new Response(rs);
630 await rejects(response.text(), TypeError);
631 },
632};
633 
634export const jsNotBytesInStart = {
635 async test() {
636 const rs = new ReadableStream({
637 start(c) {
638 c.enqueue('hello');
639 c.close();
640 },
641 });
642 
643 const response = new Response(rs);
644 await rejects(response.text(), TypeError);
645 },
646};
647 
648export const jsTeeError = {
649 async test() {
650 const rs = new ReadableStream({
651 pull(c) {
652 throw new Error('boom');
653 },
654 });
655 
656 const [branch] = rs.tee();
657 const response = new Response(branch);
658 await rejects(response.text(), { name: 'Error', message: 'boom' });
659 },
660};
661 
662export const jsTeeErrorByob = {
663 async test() {
664 const rs = new ReadableStream({
665 type: 'bytes',
666 async pull(c) {
667 if (c.byobRequest) throw new Error('boom');
668 },
669 });
670 
671 const [branch] = rs.tee();
672 const response = new Response(branch);
673 await rejects(response.text(), { name: 'Error', message: 'boom' });
674 },
675};
676 
677export const jsSourceTeed = {
678 async test() {
679 const enc = new TextEncoder();
680 
681 const readable = new ReadableStream({
682 start(c) {
683 c.enqueue(enc.encode('hello'));
684 },
685 async pull(c) {
686 c.enqueue(enc.encode('a'.repeat(100)));
687 c.enqueue(enc.encode('b'.repeat(200)));
688 c.enqueue(enc.encode('c'.repeat(300)));
689 c.enqueue(enc.encode('d'.repeat(400)));
690 c.enqueue(enc.encode('e'.repeat(500)));
691 c.enqueue(enc.encode('f'.repeat(600)));
692 c.enqueue(enc.encode('g'.repeat(700)));
693 c.enqueue(enc.encode('h'.repeat(800)));
694 await scheduler.wait(10);
695 c.enqueue(enc.encode('i'.repeat(900)));
696 c.enqueue(enc.encode('j'.repeat(1000)));
697 c.enqueue(enc.encode('k'.repeat(1100)));
698 c.close();
699 },
700 });
701 
702 const tee = readable.tee();
703 
704 async function consume(branch) {
705 for await (const _chunk of branch) {
706 // intentionally empty
707 }
708 }
709 
710 const consumePromise = consume(tee[1]);
711 
712 const response = new Response(tee[0]);
713 const text = await response.text();
714 
715 await consumePromise;
716 
717 strictEqual(text.startsWith('hello'), true);
718 strictEqual(text.endsWith('k'.repeat(1100)), true);
719 strictEqual(
720 text.length,
721 5 + 100 + 200 + 300 + 400 + 500 + 600 + 700 + 800 + 900 + 1000 + 1100
722 );
723 },
724};
725 
726export const jsByteSourceLargeData = {
727 async test() {
728 const largeData = 'x'.repeat(100000);
729 const enc = new TextEncoder();
730 
731 const sourceReadable = new ReadableStream({
732 type: 'bytes',
733 start(c) {
734 c.enqueue(enc.encode(largeData));
735 c.close();
736 },
737 });
738 
739 const reader = sourceReadable.getReader({ mode: 'byob' });
740 
741 const rs = new ReadableStream({
742 type: 'bytes',
743 async pull(c) {
744 const request = c.byobRequest;
745 if (request != null) {
746 const chunk = await reader.read(request.view);
747 if (!chunk.done) {
748 request.respondWithNewView(chunk.value);
749 } else {
750 c.close();
751 }
752 }
753 },
754 });
755 
756 const response = new Response(rs);
757 const text = await response.text();
758 strictEqual(text.length, 100000);
759 strictEqual(text, largeData);
760 },
761};
762 
763export const jsByteSourceLargeDataEnqueue = {
764 async test() {
765 const largeData = 'x'.repeat(100000);
766 const enc = new TextEncoder();
767 
768 const sourceReadable = new ReadableStream({
769 type: 'bytes',
770 start(c) {
771 c.enqueue(enc.encode(largeData));
772 c.close();
773 },
774 });
775 
776 const rs = new ReadableStream({
777 type: 'bytes',
778 async pull(c) {
779 for await (const chunk of sourceReadable) {
780 c.enqueue(chunk);
781 }
782 c.close();
783 },
784 });
785 
786 const response = new Response(rs);
787 const text = await response.text();
788 strictEqual(text.length, 100000);
789 strictEqual(text, largeData);
790 },
791};