Skip to content
File

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

javascript758 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, ok, rejects, deepStrictEqual, throws } from 'node:assert';
6import { mock } from 'node:test';
7 
8// Test pipeThrough from JavaScript readable to internal writable
9export const pipeThroughJsToInternal = {
10 async test() {
11 const enc = new TextEncoder();
12 const dec = new TextDecoder();
13 const chunks = [enc.encode('hello'), enc.encode('there'), 'hello'];
14 const rs = new ReadableStream({
15 pull(c) {
16 c.enqueue(chunks.shift());
17 if (chunks.length === 0) c.close();
18 },
19 });
20 const transform = new IdentityTransformStream();
21 const readable = rs.pipeThrough(transform);
22 
23 const output = [];
24 async function consumeStream() {
25 for await (const chunk of readable) {
26 output.push(dec.decode(chunk));
27 }
28 }
29 // The 'hello' string at the end of chunks will cause an error to be thrown.
30 await rejects(consumeStream, {
31 message: 'This WritableStream only supports writing byte types.',
32 });
33 
34 deepStrictEqual(output, ['hello', 'there']);
35 },
36};
37 
38// Test pipeThrough error in JS Readable aborts Internal Writable when preventAbort = false
39export const pipeThroughJsToInternalErroredSource = {
40 async test() {
41 const enc = new TextEncoder();
42 const rs = new ReadableStream({
43 async pull() {
44 throw new Error('boom');
45 },
46 });
47 const transform = new IdentityTransformStream();
48 const readable = rs.pipeThrough(transform);
49 
50 ok(transform.writable.locked);
51 
52 const reader = readable.getReader();
53 
54 await rejects(reader.read(), { message: 'boom' });
55 
56 ok(!transform.writable.locked);
57 
58 // Attempts to use the writable from here on will fail with the same error.
59 const writer = transform.writable.getWriter();
60 await rejects(writer.write(enc.encode('hello')), { message: 'boom' });
61 },
62};
63 
64// Test pipeTo error in JS Readable aborts Internal Writable when preventAbort = false
65export const pipeToJsToInternalErroredSource = {
66 async test() {
67 const enc = new TextEncoder();
68 const rs = new ReadableStream({
69 async pull() {
70 throw new Error('boom');
71 },
72 });
73 const { readable, writable } = new IdentityTransformStream();
74 const pipe = rs.pipeTo(writable);
75 
76 ok(writable.locked);
77 
78 const reader = readable.getReader();
79 
80 await rejects(reader.read(), { message: 'boom' });
81 
82 ok(!writable.locked);
83 
84 // Attempts to use the writable from here on will fail with the same error.
85 const writer = writable.getWriter();
86 await rejects(writer.write(enc.encode('hello')), { message: 'boom' });
87 
88 await rejects(pipe, { message: 'boom' });
89 },
90};
91 
92// Test pipeThrough error in JS Readable does not abort Internal Writable when preventAbort = true
93export const pipeThroughJsToInternalErroredSourcePreventAbort = {
94 async test() {
95 const enc = new TextEncoder();
96 const _dec = new TextDecoder();
97 const transform = new IdentityTransformStream();
98 const rs = new ReadableStream({
99 async pull() {
100 throw new Error('boom');
101 },
102 });
103 const readable = rs.pipeThrough(transform, { preventAbort: true });
104 
105 let reader = readable.getReader();
106 
107 ok(rs.locked);
108 ok(transform.writable.locked);
109 ok(transform.readable.locked);
110 
111 // Allow the piping algorithm to process the error from the pull.
112 await scheduler.wait(1);
113 
114 reader.releaseLock();
115 ok(!rs.locked);
116 ok(!transform.readable.locked);
117 ok(!transform.writable.locked);
118 
119 // We can still use the transform's readable and writable here.
120 const writer = transform.writable.getWriter();
121 reader = transform.readable.getReader();
122 
123 await Promise.all([writer.write(enc.encode('hello')), reader.read()]);
124 },
125};
126 
127// Test pipeTo error in JS Readable does not abort Internal Writable when preventAbort = true
128export const pipeToJsToInternalErroredSourcePreventAbort = {
129 async test() {
130 const enc = new TextEncoder();
131 const dec = new TextDecoder();
132 const { writable, readable } = new IdentityTransformStream();
133 const rs = new ReadableStream({
134 async pull() {
135 throw new Error('boom');
136 },
137 });
138 const pipe = rs.pipeTo(writable, { preventAbort: true });
139 
140 let reader = readable.getReader();
141 
142 ok(rs.locked);
143 ok(writable.locked);
144 ok(readable.locked);
145 
146 // The pipe promise should be rejected here but the WritableStream
147 // destination should still be usable.
148 await rejects(pipe, { message: 'boom' });
149 
150 reader.releaseLock();
151 ok(!rs.locked);
152 ok(!readable.locked);
153 ok(!writable.locked);
154 
155 // We can still use the transform's readable and writable here.
156 const writer = writable.getWriter();
157 reader = readable.getReader();
158 writer.write(enc.encode('hello'));
159 writer.close();
160 const result = await reader.read();
161 strictEqual(dec.decode(result.value), 'hello');
162 },
163};
164 
165// Test pipeThrough error in Writable cancels Readable when preventCancel = false
166export const pipeThroughJsToInternalErroredDest = {
167 async test() {
168 const enc = new TextEncoder();
169 const transform = new IdentityTransformStream();
170 const rs = new ReadableStream({
171 start(c) {
172 c.enqueue(enc.encode('hello'));
173 },
174 });
175 const readable = rs.pipeThrough(transform);
176 
177 const reader = readable.getReader();
178 
179 ok(rs.locked);
180 ok(transform.writable.locked);
181 
182 reader.cancel(new Error('boom'));
183 reader.releaseLock();
184 
185 // Allow the cancel to propagate back to the source.
186 await scheduler.wait(1);
187 
188 ok(!rs.locked);
189 ok(!transform.readable.locked);
190 ok(!transform.writable.locked);
191 
192 // Our JavaScript ReadableStream should be closed (not errored).
193 // Cancel propagates back and closes the source stream.
194 const reader2 = rs.getReader();
195 const result = await reader2.read();
196 ok(result.done);
197 strictEqual(result.value, undefined);
198 },
199};
200 
201// Test pipeTo error in Writable cancels Readable when preventCancel = false
202export const pipeToJsToInternalErroredDest = {
203 async test() {
204 const enc = new TextEncoder();
205 const { readable, writable } = new IdentityTransformStream();
206 const rs = new ReadableStream({
207 start(c) {
208 c.enqueue(enc.encode('hello'));
209 },
210 });
211 const pipe = rs.pipeTo(writable);
212 
213 const reader = readable.getReader();
214 
215 ok(rs.locked);
216 ok(writable.locked);
217 
218 reader.cancel(new Error('boom'));
219 reader.releaseLock();
220 
221 await rejects(pipe, { message: 'boom' });
222 
223 ok(!rs.locked);
224 ok(!readable.locked);
225 ok(!writable.locked);
226 
227 // Our JavaScript ReadableStream should be closed (not errored).
228 // Cancel propagates back and closes the source stream.
229 const reader2 = rs.getReader();
230 const result = await reader2.read();
231 ok(result.done);
232 strictEqual(result.value, undefined);
233 },
234};
235 
236// Test closing Readable closes Writable when preventClose = false
237export const pipeThroughJsToInternalCloses = {
238 async test() {
239 const enc = new TextEncoder();
240 const chunks = [enc.encode('hello'), enc.encode('there')];
241 const rs = new ReadableStream({
242 pull(c) {
243 c.enqueue(chunks.shift());
244 if (chunks.length === 0) c.close();
245 },
246 });
247 const transform = new IdentityTransformStream();
248 const readable = rs.pipeThrough(transform);
249 
250 for await (const _chunk of readable) {
251 // consume all chunks
252 }
253 
254 // The writable should be closed and locked by the pipe
255 throws(() => transform.writable.getWriter(), {
256 message: 'This WritableStream is currently locked to a writer.',
257 });
258 },
259};
260 
261// Test closing Readable does not close Writable when preventClose = true
262export const pipeThroughJsToInternalPreventClose = {
263 async test() {
264 const enc = new TextEncoder();
265 const dec = new TextDecoder();
266 const rs = new ReadableStream({
267 start(c) {
268 c.enqueue(enc.encode('hello'));
269 c.close();
270 },
271 });
272 const transform = new IdentityTransformStream();
273 const readable = rs.pipeThrough(transform, { preventClose: true });
274 
275 // Because the internal TransformStream here won't resolve the write
276 // promises until a read has been performed, we have to read, then
277 // wait a turn of the event loop before we can check that the writer
278 // is still in the correct state.
279 const reader = readable.getReader();
280 await reader.read();
281 
282 await scheduler.wait(1);
283 
284 // The WritableStream should not be closed and still usable.
285 const writer = transform.writable.getWriter();
286 writer.write(enc.encode('there'));
287 const read = await reader.read();
288 strictEqual(dec.decode(read.value), 'there');
289 },
290};
291 
292// Test pipeThrough with BYOB ReadableStream works
293export const pipeThroughJsByobToInternal = {
294 async test() {
295 const enc = new TextEncoder();
296 const dec = new TextDecoder();
297 const chunks = [enc.encode('hello'), enc.encode('there')];
298 const rs = new ReadableStream({
299 type: 'bytes',
300 pull(c) {
301 c.enqueue(chunks.shift());
302 if (chunks.length === 0) c.close();
303 },
304 });
305 const transform = new IdentityTransformStream();
306 const readable = rs.pipeThrough(transform);
307 
308 const output = [];
309 for await (const chunk of readable) {
310 output.push(dec.decode(chunk));
311 }
312 
313 strictEqual(output[0], 'hello');
314 strictEqual(output[1], 'there');
315 },
316};
317 
318// Test pipeTo with BYOB ReadableStream works
319export const pipeToJsByobToInternal = {
320 async test() {
321 const enc = new TextEncoder();
322 const dec = new TextDecoder();
323 const chunks = [enc.encode('hello'), enc.encode('there')];
324 const rs = new ReadableStream({
325 type: 'bytes',
326 pull(c) {
327 c.enqueue(chunks.shift());
328 if (chunks.length === 0) c.close();
329 },
330 });
331 const { readable, writable } = new IdentityTransformStream();
332 rs.pipeTo(writable);
333 
334 const output = [];
335 for await (const chunk of readable) {
336 output.push(dec.decode(chunk));
337 }
338 
339 strictEqual(output[0], 'hello');
340 strictEqual(output[1], 'there');
341 },
342};
343 
344// Test simple pipeTo from internal readable to JavaScript writable
345export const pipeToInternalToJsSimple = {
346 async test() {
347 const enc = new TextEncoder();
348 const dec = new TextDecoder();
349 
350 const { readable, writable } = new IdentityTransformStream();
351 
352 const chunks = [];
353 const ws = new WritableStream({
354 write(chunk) {
355 chunks.push(chunk);
356 },
357 });
358 
359 const pipe = readable.pipeTo(ws);
360 
361 const writer = writable.getWriter();
362 writer.write(enc.encode('hello'));
363 writer.write(enc.encode('there'));
364 writer.close();
365 
366 await pipe;
367 
368 strictEqual(dec.decode(chunks[0]), 'hello');
369 strictEqual(dec.decode(chunks[1]), 'there');
370 
371 ok(!ws.locked);
372 ok(!readable.locked);
373 
374 const writer2 = ws.getWriter();
375 await rejects(writer2.write('no'), {
376 message: 'This WritableStream has been closed.',
377 });
378 },
379};
380 
381// Test pipeTo error in internal readable aborts JS writable when preventAbort = false
382export const pipeToInternalToJsError = {
383 async test() {
384 const _enc = new TextEncoder();
385 
386 const { readable, writable } = new IdentityTransformStream();
387 
388 const ws = new WritableStream({
389 write(chunk) {},
390 });
391 
392 const pipe = readable.pipeTo(ws);
393 
394 writable.abort(new Error('boom'));
395 
396 await rejects(pipe, { message: 'boom' });
397 
398 const writer = ws.getWriter();
399 await rejects(writer.write('hello'), { message: 'boom' });
400 },
401};
402 
403// Test pipeTo error in internal readable does not abort JS writable when preventAbort = true
404export const pipeToInternalToJsErrorPrevent = {
405 async test() {
406 const { readable, writable } = new IdentityTransformStream();
407 
408 const ws = new WritableStream({
409 write(chunk) {},
410 });
411 
412 const pipe = readable.pipeTo(ws, { preventAbort: true });
413 
414 writable.abort(new Error('boom'));
415 
416 await rejects(pipe, { message: 'boom' });
417 
418 ok(!ws.locked);
419 
420 const writer = ws.getWriter();
421 await writer.write('hello');
422 },
423};
424 
425// Test pipeTo closing internal readable closes JS writable when preventClose = false
426export const pipeToInternalToJsClose = {
427 async test() {
428 const { readable, writable } = new IdentityTransformStream();
429 
430 const ws = new WritableStream({});
431 
432 readable.pipeTo(ws);
433 
434 writable.close();
435 
436 // Allow the close to propagate through the pipe.
437 await scheduler.wait(1);
438 
439 const writer = ws.getWriter();
440 await rejects(writer.write('hello'), {
441 message: 'This WritableStream has been closed.',
442 });
443 },
444};
445 
446// Test pipeTo closing internal readable does not close JS writable when preventClose = true
447export const pipeToInternalToJsClosePrevent = {
448 async test() {
449 const { readable, writable } = new IdentityTransformStream();
450 
451 const ws = new WritableStream({});
452 
453 readable.pipeTo(ws, { preventClose: true });
454 
455 writable.close();
456 
457 // Allow the pipe to finish without closing the destination.
458 await scheduler.wait(1);
459 
460 const writer = ws.getWriter();
461 await writer.write('hello');
462 },
463};
464 
465// Test simple pipeTo JS-to-JS
466export const pipeToJsToJsSimple = {
467 async test() {
468 const chunks = [1, 2, 3];
469 const output = [];
470 const readable = new ReadableStream({
471 async pull(c) {
472 c.enqueue(chunks.shift());
473 if (chunks.length === 0) c.close();
474 },
475 });
476 const writable = new WritableStream({
477 write(chunk) {
478 output.push(chunk);
479 },
480 });
481 
482 await readable.pipeTo(writable);
483 
484 deepStrictEqual(output, [1, 2, 3]);
485 },
486};
487 
488// Test pipeTo error in JS readable aborts JS writable when preventAbort = false
489export const pipeToJsToJsErrorReadable = {
490 async test() {
491 const abortFn = mock.fn();
492 const readable = new ReadableStream({
493 async pull() {
494 throw new Error('boom');
495 },
496 });
497 const writable = new WritableStream({
498 abort: abortFn,
499 });
500 
501 await rejects(readable.pipeTo(writable), { message: 'boom' });
502 
503 strictEqual(abortFn.mock.callCount(), 1);
504 },
505};
506 
507// Test pipeTo error in JS readable does not abort JS writable when preventAbort = true
508export const pipeToJsToJsErrorReadablePrevent = {
509 async test() {
510 const abortFn = mock.fn();
511 const readable = new ReadableStream({
512 async pull() {
513 throw new Error('boom');
514 },
515 });
516 const writable = new WritableStream({
517 abort: abortFn,
518 });
519 
520 const pipe = readable.pipeTo(writable, { preventAbort: true });
521 
522 await rejects(pipe, { message: 'boom' });
523 
524 strictEqual(abortFn.mock.callCount(), 0);
525 },
526};
527 
528// Test pipeTo error in JS writable cancels JS readable when preventCancel = false
529export const pipeToJsToJsErrorWritable = {
530 async test() {
531 const cancelFn = mock.fn();
532 const readable = new ReadableStream({
533 start(c) {
534 c.enqueue('hello');
535 },
536 cancel: cancelFn,
537 });
538 const writable = new WritableStream({
539 write() {
540 throw new Error('boom');
541 },
542 });
543 
544 const pipe = readable.pipeTo(writable);
545 
546 await rejects(pipe, { message: 'boom' });
547 
548 strictEqual(cancelFn.mock.callCount(), 1);
549 },
550};
551 
552// Test pipeTo error in JS writable does not cancel JS readable when preventCancel = true
553export const pipeToJsToJsErrorWritablePrevent = {
554 async test() {
555 const chunks = [1, 2];
556 const cancelFn = mock.fn();
557 const readable = new ReadableStream({
558 pull(c) {
559 c.enqueue(chunks.shift());
560 if (chunks.length === 0) c.close();
561 },
562 cancel: cancelFn,
563 });
564 const writable = new WritableStream({
565 write() {
566 throw new Error('boom');
567 },
568 });
569 
570 const pipe = readable.pipeTo(writable, { preventCancel: true });
571 
572 await rejects(pipe, { message: 'boom' });
573 
574 strictEqual(cancelFn.mock.callCount(), 0);
575 
576 const reader = readable.getReader();
577 await reader.read();
578 },
579};
580 
581// Test closing JS readable closes JS writable when preventClose = false
582export const pipeToJsToJsCloseReadable = {
583 async test() {
584 const closeFn = mock.fn();
585 const readable = new ReadableStream({
586 start(c) {
587 c.close();
588 },
589 });
590 const writable = new WritableStream({
591 close: closeFn,
592 });
593 
594 await readable.pipeTo(writable);
595 
596 strictEqual(closeFn.mock.callCount(), 1);
597 },
598};
599 
600// Test closing JS readable does not close JS writable when preventClose = true
601export const pipeToJsToJsCloseReadablePrevent = {
602 async test() {
603 const closeFn = mock.fn();
604 const readable = new ReadableStream({
605 start(c) {
606 c.close();
607 },
608 });
609 const writable = new WritableStream({
610 close: closeFn,
611 });
612 
613 await readable.pipeTo(writable, { preventClose: true });
614 
615 strictEqual(closeFn.mock.callCount(), 0);
616 },
617};
618 
619// Test pipeTo from a tee branch
620export const pipeToJsToJsTee = {
621 async test() {
622 const readable = new ReadableStream({
623 start(c) {
624 c.enqueue('hello');
625 c.close();
626 },
627 });
628 
629 const output = [];
630 const writable = new WritableStream({
631 write(chunk) {
632 output.push(chunk);
633 },
634 });
635 
636 const [branch] = readable.tee();
637 
638 await branch.pipeTo(writable);
639 
640 strictEqual(output[0], 'hello');
641 },
642};
643 
644// Test pipeTo with already-aborted signal (JS to JS)
645export const pipeToJsToJsCancelAlready = {
646 async test() {
647 const signal = AbortSignal.abort(new Error('boom'));
648 const writeFn = mock.fn();
649 const readable = new ReadableStream({
650 start(c) {
651 c.enqueue('hello');
652 c.close();
653 },
654 });
655 
656 const writable = new WritableStream({
657 write: writeFn,
658 });
659 
660 await rejects(readable.pipeTo(writable, { signal }), { message: 'boom' });
661 
662 strictEqual(writeFn.mock.callCount(), 0);
663 },
664};
665 
666// Test pipeTo with already-aborted signal (JS to native)
667export const pipeToJsToNativeCancelAlready = {
668 async test() {
669 const signal = AbortSignal.abort(new Error('boom'));
670 const source = new ReadableStream({
671 start(c) {
672 c.enqueue('hello');
673 c.close();
674 },
675 });
676 
677 const { writable, readable: _readable } = new TransformStream();
678 
679 await rejects(source.pipeTo(writable, { signal }), { message: 'boom' });
680 },
681};
682 
683// Test pipeTo cancelable during operation (JS to JS)
684export const pipeToJsToJsCancel = {
685 async test() {
686 const controller = new AbortController();
687 
688 const readable = new ReadableStream({
689 start(c) {
690 c.enqueue('hello');
691 },
692 });
693 
694 let output = '';
695 const writable = new WritableStream({
696 write(chunk) {
697 output += chunk;
698 controller.abort(new Error('boom'));
699 },
700 });
701 
702 await rejects(readable.pipeTo(writable, { signal: controller.signal }), {
703 message: 'boom',
704 });
705 
706 strictEqual(output, 'hello');
707 },
708};
709 
710// Test pipeTo cancelable during operation (JS to native)
711export const pipeToJsToNativeCancel = {
712 async test() {
713 const controller = new AbortController();
714 const enc = new TextEncoder();
715 
716 let ready = false;
717 
718 const source = new ReadableStream({
719 start(c) {
720 c.enqueue(enc.encode('hello'));
721 },
722 pull() {
723 if (ready) {
724 controller.abort(new Error('boom'));
725 }
726 },
727 });
728 
729 const { readable, writable } = new TransformStream();
730 
731 const reader = readable.getReader();
732 
733 ready = true;
734 
735 const promises = await Promise.allSettled([
736 source.pipeTo(writable, { signal: controller.signal }),
737 reader.read(),
738 ]);
739 
740 strictEqual(promises[0].status, 'rejected');
741 strictEqual(promises[1].status, 'rejected');
742 
743 strictEqual(promises[0].reason.message, 'boom');
744 strictEqual(promises[1].reason.message, 'boom');
745 },
746};
747 
748// Default fetch handler for service binding requests
749export default {
750 async fetch(request) {
751 if (request.url.includes('/stream')) {
752 const data = 'hello world '.repeat(100);
753 return new Response(data);
754 }
755 return new Response('Not found', { status: 404 });
756 },
757};