Skip to content
File

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

javascript674 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 } from 'node:assert';
6 
7async function consume(readable) {
8 let data = '';
9 for await (const chunk of readable) {
10 data += chunk;
11 }
12 return data;
13}
14 
15// Test default identity transform
16export const defaultIdentityTransform = {
17 async test() {
18 const transform = new TransformStream();
19 ok(!(transform instanceof IdentityTransformStream));
20 
21 const { readable, writable } = transform;
22 
23 const writer = writable.getWriter();
24 const reader = readable.getReader();
25 
26 const after = await Promise.allSettled([
27 writer.write('hello'),
28 reader.read(),
29 ]);
30 
31 strictEqual(after[0].value, undefined);
32 strictEqual(after[1].value.value, 'hello');
33 },
34};
35 
36// Test simple transform with start/transform/flush
37export const simpleTransform = {
38 async test() {
39 const transform = new TransformStream({
40 start(controller) {
41 controller.enqueue('<');
42 },
43 transform(value, controller) {
44 controller.enqueue(value.toUpperCase());
45 },
46 flush(controller) {
47 controller.enqueue('>');
48 },
49 });
50 ok(!(transform instanceof IdentityTransformStream));
51 
52 const { readable, writable } = transform;
53 
54 const writer = writable.getWriter();
55 
56 const res = await Promise.allSettled([
57 writer.write('hello'),
58 writer.write('there'),
59 writer.close(),
60 consume(readable),
61 ]);
62 
63 strictEqual(res[3].value, '<HELLOTHERE>');
64 },
65};
66 
67// Test async transform with delays
68export const delayTransform = {
69 async test() {
70 const transform = new TransformStream({
71 async start(controller) {
72 await scheduler.wait(1);
73 controller.enqueue('<');
74 },
75 async transform(value, controller) {
76 await scheduler.wait(1);
77 controller.enqueue(value.toUpperCase());
78 },
79 async flush(controller) {
80 await scheduler.wait(1);
81 controller.enqueue('>');
82 },
83 });
84 ok(!(transform instanceof IdentityTransformStream));
85 
86 const { readable, writable } = transform;
87 
88 const writer = writable.getWriter();
89 
90 const res = await Promise.allSettled([
91 writer.write('hello'),
92 writer.write('there'),
93 writer.close(),
94 consume(readable),
95 ]);
96 
97 strictEqual(res[3].value, '<HELLOTHERE>');
98 },
99};
100 
101// Test transform with different types (string input, Uint8Array output)
102export const differentTypesTransform = {
103 async test() {
104 const enc = new TextEncoder();
105 
106 async function consumeBytes(readable) {
107 const dec = new TextDecoder();
108 let data = '';
109 for await (const chunk of readable) {
110 data += dec.decode(chunk, { stream: true });
111 }
112 data += dec.decode();
113 return data;
114 }
115 
116 const transform = new TransformStream({
117 async start(controller) {
118 await scheduler.wait(1);
119 controller.enqueue(enc.encode('<'));
120 },
121 async transform(value, controller) {
122 await scheduler.wait(1);
123 controller.enqueue(enc.encode(value.toUpperCase()));
124 },
125 async flush(controller) {
126 await scheduler.wait(1);
127 controller.enqueue(enc.encode('>'));
128 },
129 });
130 ok(!(transform instanceof IdentityTransformStream));
131 
132 const { readable, writable } = transform;
133 
134 const writer = writable.getWriter();
135 
136 const res = await Promise.allSettled([
137 writer.write('hello'),
138 writer.write('there'),
139 writer.close(),
140 consumeBytes(readable),
141 ]);
142 
143 strictEqual(res[3].value, '<HELLOTHERE>');
144 },
145};
146 
147// Test sync error during start
148export const syncErrorDuringStart = {
149 async test() {
150 const transform = new TransformStream({
151 start() {
152 throw new Error('boom');
153 },
154 });
155 ok(!(transform instanceof IdentityTransformStream));
156 
157 const { readable, writable } = transform;
158 
159 const writer = writable.getWriter();
160 
161 const res = await Promise.allSettled([
162 writer.write('hello'),
163 writer.write('there'),
164 writer.close(),
165 consume(readable),
166 ]);
167 
168 strictEqual(res[0].status, 'rejected');
169 strictEqual(res[1].status, 'rejected');
170 strictEqual(res[2].status, 'rejected');
171 strictEqual(res[0].reason.message, 'boom');
172 strictEqual(res[1].reason.message, 'boom');
173 strictEqual(res[2].reason.message, 'boom');
174 },
175};
176 
177// Test async error during start
178export const asyncErrorDuringStart = {
179 async test() {
180 const transform = new TransformStream({
181 async start() {
182 throw new Error('boom');
183 },
184 });
185 ok(!(transform instanceof IdentityTransformStream));
186 
187 const { readable, writable } = transform;
188 
189 const writer = writable.getWriter();
190 
191 const res = await Promise.allSettled([
192 writer.write('hello'),
193 writer.write('there'),
194 writer.close(),
195 consume(readable),
196 ]);
197 
198 strictEqual(res[0].status, 'rejected');
199 strictEqual(res[1].status, 'rejected');
200 strictEqual(res[2].status, 'rejected');
201 strictEqual(res[3].status, 'rejected');
202 strictEqual(res[0].reason.message, 'boom');
203 strictEqual(res[1].reason.message, 'boom');
204 strictEqual(res[2].reason.message, 'boom');
205 strictEqual(res[3].reason.message, 'boom');
206 },
207};
208 
209// Test sync error during transform
210export const syncErrorDuringTransform = {
211 async test() {
212 const transform = new TransformStream({
213 transform() {
214 throw new Error('boom');
215 },
216 });
217 ok(!(transform instanceof IdentityTransformStream));
218 
219 const { readable, writable } = transform;
220 
221 const writer = writable.getWriter();
222 
223 const res = await Promise.allSettled([
224 writer.write('hello'),
225 writer.write('there'),
226 writer.close(),
227 consume(readable),
228 ]);
229 
230 strictEqual(res[0].status, 'rejected');
231 strictEqual(res[1].status, 'rejected');
232 strictEqual(res[2].status, 'rejected');
233 strictEqual(res[3].status, 'rejected');
234 strictEqual(res[0].reason.message, 'boom');
235 strictEqual(res[1].reason.message, 'boom');
236 strictEqual(res[2].reason.message, 'boom');
237 strictEqual(res[3].reason.message, 'boom');
238 },
239};
240 
241// Test async error during transform
242export const asyncErrorDuringTransform = {
243 async test() {
244 const transform = new TransformStream({
245 async transform() {
246 await scheduler.wait(1);
247 throw new Error('boom');
248 },
249 });
250 ok(!(transform instanceof IdentityTransformStream));
251 
252 const { readable, writable } = transform;
253 
254 const writer = writable.getWriter();
255 
256 const res = await Promise.allSettled([
257 writer.write('hello'),
258 writer.write('there'),
259 writer.close(),
260 consume(readable),
261 ]);
262 
263 strictEqual(res[0].status, 'rejected');
264 strictEqual(res[1].status, 'rejected');
265 strictEqual(res[2].status, 'rejected');
266 strictEqual(res[3].status, 'rejected');
267 strictEqual(res[0].reason.message, 'boom');
268 strictEqual(res[1].reason.message, 'boom');
269 strictEqual(res[2].reason.message, 'boom');
270 strictEqual(res[3].reason.message, 'boom');
271 },
272};
273 
274// Test sync error during flush
275export const syncErrorDuringFlush = {
276 async test() {
277 const transform = new TransformStream({
278 flush() {
279 throw new Error('boom');
280 },
281 });
282 ok(!(transform instanceof IdentityTransformStream));
283 
284 const { readable, writable } = transform;
285 
286 const writer = writable.getWriter();
287 
288 const res = await Promise.allSettled([
289 writer.write('hello'),
290 writer.write('there'),
291 writer.close(),
292 consume(readable),
293 ]);
294 
295 // The two writes will succeed.
296 strictEqual(res[0].status, 'fulfilled');
297 strictEqual(res[1].status, 'fulfilled');
298 
299 // The close and the consume will reject.
300 strictEqual(res[2].status, 'rejected');
301 strictEqual(res[3].status, 'rejected');
302 strictEqual(res[2].reason.message, 'boom');
303 strictEqual(res[3].reason.message, 'boom');
304 },
305};
306 
307// Test async error during flush
308export const asyncErrorDuringFlush = {
309 async test() {
310 const transform = new TransformStream({
311 async flush() {
312 await scheduler.wait(1);
313 throw new Error('boom');
314 },
315 });
316 ok(!(transform instanceof IdentityTransformStream));
317 
318 const { readable, writable } = transform;
319 
320 const writer = writable.getWriter();
321 
322 const res = await Promise.allSettled([
323 writer.write('hello'),
324 writer.write('there'),
325 writer.close(),
326 consume(readable),
327 ]);
328 
329 // The two writes will succeed.
330 strictEqual(res[0].status, 'fulfilled');
331 strictEqual(res[1].status, 'fulfilled');
332 
333 // The close and the consume will reject.
334 strictEqual(res[2].status, 'rejected');
335 strictEqual(res[3].status, 'rejected');
336 strictEqual(res[2].reason.message, 'boom');
337 strictEqual(res[3].reason.message, 'boom');
338 },
339};
340 
341// Test write backpressure
342export const writeBackpressure = {
343 async test() {
344 let expectedReadSize = 2;
345 const transform = new TransformStream(
346 {
347 transform(chunk, controller) {
348 strictEqual(controller.desiredSize, expectedReadSize--);
349 controller.enqueue(chunk);
350 },
351 },
352 { highWaterMark: 2 },
353 { highWaterMark: 2 }
354 );
355 
356 const writer = transform.writable.getWriter();
357 strictEqual(writer.desiredSize, 2);
358 
359 const promises = [writer.write('hello'), writer.write('there')];
360 
361 strictEqual(writer.desiredSize, 0);
362 
363 await Promise.allSettled(promises);
364 
365 strictEqual(writer.desiredSize, 2);
366 },
367};
368 
369// Test that piping from a JS-backed TransformStream through an
370// IdentityTransformStream does not result in a hung pipeTo promise.
371export const transformRoundtrip = {
372 async test(ctrl, env, ctx) {
373 const enc = new TextEncoder();
374 const dec = new TextDecoder();
375 
376 const testData = 'hello world test data';
377 const compressedStream = new ReadableStream({
378 start(controller) {
379 controller.enqueue(enc.encode(testData));
380 controller.close();
381 },
382 }).pipeThrough(new CompressionStream('gzip'));
383 
384 const compressedChunks = [];
385 const compressedReader = compressedStream.getReader();
386 for (;;) {
387 const { done, value } = await compressedReader.read();
388 if (done) break;
389 compressedChunks.push(value);
390 }
391 const compressedData = new Uint8Array(
392 compressedChunks.reduce((acc, chunk) => acc + chunk.length, 0)
393 );
394 let offset = 0;
395 for (const chunk of compressedChunks) {
396 compressedData.set(chunk, offset);
397 offset += chunk.length;
398 }
399 
400 const inputStream = new ReadableStream({
401 start(controller) {
402 controller.enqueue(compressedData);
403 controller.close();
404 },
405 });
406 
407 const decompression = new DecompressionStream('gzip');
408 const ts = new TransformStream({
409 transform(chunk, controller) {
410 controller.enqueue(chunk);
411 },
412 });
413 const { readable, writable } = new IdentityTransformStream();
414 
415 ctx.waitUntil(
416 inputStream.pipeThrough(decompression).pipeThrough(ts).pipeTo(writable)
417 );
418 
419 const outputChunks = [];
420 const reader = readable.getReader();
421 for (;;) {
422 const { done, value } = await reader.read();
423 if (done) break;
424 outputChunks.push(value);
425 }
426 
427 const output = dec.decode(
428 new Uint8Array(
429 outputChunks.reduce((acc, chunk) => [...acc, ...chunk], [])
430 )
431 );
432 strictEqual(output, testData);
433 },
434};
435 
436// Regression test: calling controller.error() from inside a readableStrategy size()
437// callback during TransformStream enqueue must not cause a use-after-free.
438// The size() callback is invoked by ReadableStreamDefaultController::enqueue().
439// Calling transformController.error() from within it drops all jsg::Ref references
440// to the ReadableStreamDefaultController, which would free it while its enqueue()
441// method is still on the stack.
442export const sizeCallbackErrorDoesNotUAF = {
443 async test() {
444 let transformCtrl;
445 const ts = new TransformStream(
446 {
447 start(controller) {
448 transformCtrl = controller;
449 },
450 transform(chunk, controller) {
451 controller.enqueue(chunk);
452 },
453 },
454 {},
455 {
456 size(_chunk) {
457 // Calling error() here drops the ReadableStreamDefaultController refs.
458 // Without the fix this is a use-after-free.
459 transformCtrl.error(new Error('errored from size'));
460 return 1;
461 },
462 highWaterMark: 1,
463 }
464 );
465 
466 const writer = ts.writable.getWriter();
467 const reader = ts.readable.getReader();
468 
469 // The write triggers transform -> enqueue -> size() -> error().
470 // The key assertion: we reach this point without crashing (no UAF).
471 const results = await Promise.allSettled([
472 writer.write('hello'),
473 reader.read(),
474 ]);
475 
476 // The read must reject because the readable side was errored.
477 strictEqual(results[1].status, 'rejected');
478 ok(results[1].reason.message.includes('errored from size'));
479 },
480};
481 
482// Same as above, but the size callback throws after calling error().
483export const sizeCallbackErrorAndThrowDoesNotUAF = {
484 async test() {
485 let transformCtrl;
486 const ts = new TransformStream(
487 {
488 start(controller) {
489 transformCtrl = controller;
490 },
491 transform(chunk, controller) {
492 controller.enqueue(chunk);
493 },
494 },
495 {},
496 {
497 size(_chunk) {
498 transformCtrl.error(new Error('errored from size'));
499 throw new Error('size threw');
500 },
501 highWaterMark: 1,
502 }
503 );
504 
505 const writer = ts.writable.getWriter();
506 const reader = ts.readable.getReader();
507 
508 // The key assertion: we reach this point without crashing (no UAF).
509 const results = await Promise.allSettled([
510 writer.write('hello'),
511 reader.read(),
512 ]);
513 
514 // The read must reject because the readable side was errored.
515 strictEqual(results[1].status, 'rejected');
516 },
517};
518 
519// Regression test: iterating over globalThis properties should not crash.
520// This tests that constructing and inspecting various global objects works correctly.
521//
522// Note: We skip `process` in the iteration because when nodejs_compat is enabled,
523// the iteration logic may call process.exit() which terminates the test.
524export const transformCrashRegression = {
525 test() {
526 function iterate(obj, depth = 0, results = {}, originalKey) {
527 for (const key in obj) {
528 if (depth > 100) return results;
529 if (
530 key === 'parent' ||
531 key === 'globalThis' ||
532 key === 'self' ||
533 key === 'ServiceWorkerGlobalScope' ||
534 key === 'global' ||
535 key === 'process'
536 )
537 continue;
538 if (typeof obj[key] === 'object') {
539 results[key] = iterate(obj[key], ++depth, results[key] ?? {}, key);
540 } else {
541 const properties = new Set([
542 ...Object.getOwnPropertyNames(obj[key]),
543 ...Object.getOwnPropertySymbols(obj[key]),
544 ...Object.keys(obj[key]),
545 ]);
546 for (const prop in obj[key]) {
547 properties.add(prop);
548 }
549 if (properties.size > 0) {
550 results[key] = {
551 __proto: typeof obj[key].__proto__,
552 };
553 for (const property of properties) {
554 if (
555 property === 'caller' ||
556 property === 'callee' ||
557 property === 'arguments' ||
558 property === 'constructor'
559 )
560 continue;
561 let writeProp =
562 property === 'prototype' ? '_prototype' : property;
563 try {
564 if (typeof obj[key][property] === 'object') {
565 results[key][writeProp] = iterate(
566 obj[key][property],
567 ++depth,
568 results[key][property] ?? {},
569 key
570 );
571 } else {
572 results[key][writeProp] = typeof obj[key][property];
573 }
574 } catch (_err) {
575 // expected
576 }
577 let instance;
578 try {
579 if (property === 'prototype') {
580 try {
581 instance = new obj[key]();
582 } catch (internalConstructError) {
583 if (
584 internalConstructError.message.includes(
585 'Failed to construct'
586 ) ||
587 internalConstructError.message.includes(
588 'Illegal constructor'
589 )
590 ) {
591 instance = new obj[key]('');
592 } else {
593 throw internalConstructError;
594 }
595 }
596 } else {
597 instance = new obj[key][property]();
598 }
599 const instanceProperties = new Set([
600 ...Object.getOwnPropertyNames(instance),
601 ...Object.getOwnPropertySymbols(instance),
602 ...Object.keys(instance),
603 ...(instance.constructor
604 ? Object.getOwnPropertyNames(instance.constructor)
605 : []),
606 ...(instance.constructor
607 ? Object.getOwnPropertySymbols(instance.constructor)
608 : []),
609 ...(instance.constructor
610 ? Object.keys(instance.constructor)
611 : []),
612 ]);
613 for (const prop in instance) {
614 instanceProperties.add(prop);
615 }
616 if (instanceProperties.size > 0) {
617 if (results[key][writeProp] === undefined) {
618 results[key][writeProp] = {};
619 }
620 for (const instanceProperty of instanceProperties) {
621 try {
622 if (
623 results[key][writeProp][instanceProperty] === undefined
624 ) {
625 results[key][writeProp][instanceProperty] =
626 typeof instance[instanceProperty];
627 }
628 } catch {
629 // intentionally empty
630 }
631 }
632 } else {
633 if (results[key][writeProp] === undefined) {
634 results[key][writeProp] = typeof instance;
635 }
636 }
637 } catch (_e) {
638 if (results[key][writeProp] === undefined) {
639 results[key][writeProp] = typeof obj[key][property];
640 }
641 continue;
642 }
643 }
644 } else {
645 results[key] = typeof obj[key];
646 }
647 }
648 }
649 return results;
650 }
651 
652 function order(obj) {
653 if (typeof obj === 'object') {
654 const ordered = {};
655 Object.keys(obj)
656 .sort()
657 .forEach(function (key) {
658 ordered[key] = order(obj[key]);
659 });
660 return ordered;
661 } else {
662 return obj;
663 }
664 }
665 
666 const iterated = iterate(globalThis);
667 const ordered = order(iterated);
668 const toWrite = JSON.stringify(ordered);
669 
670 // Test passes if we get here without crashing
671 ok(toWrite.length > 0);
672 },
673};