Skip to content
File

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

javascript804 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, ok } from 'node:assert';
6 
7// Test CompressionStream/DecompressionStream with gzip
8export const compressionGzip = {
9 async test() {
10 let output;
11 const testData = 'hello'.repeat(100);
12 const enc = new TextEncoder();
13 const dec = new TextDecoder();
14 
15 {
16 const { writable, readable } = new CompressionStream('gzip');
17 
18 const writer = writable.getWriter();
19 const reader = readable.getReader();
20 
21 await writer.write(enc.encode(testData));
22 await writer.close();
23 
24 output = await reader.read();
25 }
26 
27 {
28 const { writable, readable } = new DecompressionStream('gzip');
29 
30 const writer = writable.getWriter();
31 const reader = readable.getReader();
32 
33 await writer.write(output.value);
34 await writer.close();
35 
36 const result = await reader.read();
37 
38 strictEqual(dec.decode(result.value), testData);
39 }
40 },
41};
42 
43// Test CompressionStream/DecompressionStream with deflate
44export const compressionDeflate = {
45 async test() {
46 let output;
47 const testData = 'hello'.repeat(100);
48 const enc = new TextEncoder();
49 const dec = new TextDecoder();
50 
51 {
52 const { writable, readable } = new CompressionStream('deflate');
53 
54 const writer = writable.getWriter();
55 const reader = readable.getReader();
56 
57 await writer.write(enc.encode(testData));
58 await writer.close();
59 
60 output = await reader.read();
61 }
62 
63 {
64 const { writable, readable } = new DecompressionStream('deflate');
65 
66 const writer = writable.getWriter();
67 const reader = readable.getReader();
68 
69 await writer.write(output.value);
70 await writer.close();
71 
72 const result = await reader.read();
73 
74 strictEqual(dec.decode(result.value), testData);
75 }
76 },
77};
78 
79// Test CompressionStream/DecompressionStream with deflate-raw
80export const compressionDeflateRaw = {
81 async test() {
82 let output;
83 const testData = 'hello'.repeat(100);
84 const enc = new TextEncoder();
85 const dec = new TextDecoder();
86 
87 {
88 const { writable, readable } = new CompressionStream('deflate-raw');
89 
90 const writer = writable.getWriter();
91 const reader = readable.getReader();
92 
93 await writer.write(enc.encode(testData));
94 await writer.close();
95 
96 output = await reader.read();
97 }
98 
99 {
100 const { writable, readable } = new DecompressionStream('deflate-raw');
101 
102 const writer = writable.getWriter();
103 const reader = readable.getReader();
104 
105 await writer.write(output.value);
106 await writer.close();
107 
108 const result = await reader.read();
109 
110 strictEqual(dec.decode(result.value), testData);
111 }
112 },
113};
114 
115// Test compression/decompression with pending read
116export const compressionPendingRead = {
117 async test() {
118 const testData = 'hello';
119 const check = new Uint8Array([
120 0x78, 0x9c, 0xcb, 0x48, 0xcd, 0xc9, 0xc9, 0x07, 0x00, 0x06, 0x2c, 0x02,
121 0x15,
122 ]);
123 const enc = new TextEncoder();
124 const dec = new TextDecoder();
125 
126 {
127 const { writable, readable } = new CompressionStream('deflate');
128 
129 const writer = writable.getWriter();
130 const reader = readable.getReader();
131 
132 const read = reader.read();
133 
134 await writer.write(enc.encode(testData));
135 await writer.close();
136 
137 await read;
138 }
139 
140 {
141 const { writable, readable } = new DecompressionStream('deflate');
142 
143 const writer = writable.getWriter();
144 const reader = readable.getReader();
145 
146 const read = reader.read();
147 
148 await writer.write(check);
149 await writer.close();
150 
151 const result = await read;
152 
153 strictEqual(dec.decode(result.value), testData);
154 }
155 },
156};
157 
158// Test cancel/abort behavior
159export const compressionCancelAbort = {
160 async test() {
161 {
162 const { writable, readable } = new CompressionStream('deflate');
163 const reader = readable.getReader();
164 const writer = writable.getWriter();
165 const promise = reader.read();
166 writer.abort(new Error('boom'));
167 await rejects(promise, { message: 'boom' });
168 }
169 
170 {
171 const { readable } = new CompressionStream('deflate');
172 const reader = readable.getReader();
173 const promise = reader.read();
174 reader.cancel(new Error('boom'));
175 await rejects(promise, { message: 'boom' });
176 }
177 },
178};
179 
180// Test decompression error handling
181export const decompressionError = {
182 async test() {
183 // The data is not compressed so DecompressionStream will fail.
184 const { writable, readable } = new DecompressionStream('deflate');
185 
186 const writer = writable.getWriter();
187 const reader = readable.getReader();
188 
189 // Write uncompressed data which should cause a decompression error
190 // The write itself may also reject, so we need to handle that too
191 const writePromise = writer
192 .write(new TextEncoder().encode('not compressed data'))
193 .catch(() => {});
194 
195 // The read should fail with a TypeError
196 await rejects(reader.read(), TypeError);
197 
198 // A second attempt to read also fails
199 await rejects(reader.read(), TypeError);
200 
201 // Ensure the write operation completes before the test ends
202 await writePromise;
203 },
204};
205 
206// Test decompression error via iteration
207export const decompressionErrorIteration = {
208 async test() {
209 const { writable, readable } = new DecompressionStream('deflate');
210 
211 const writer = writable.getWriter();
212 
213 // Write uncompressed data which should cause a decompression error
214 writer.write(new TextEncoder().encode('not compressed data'));
215 
216 const consume = async () => {
217 for await (const _ of readable) {
218 // Should not reach here
219 }
220 };
221 
222 await rejects(consume(), TypeError);
223 },
224};
225 
226// Test piped decompression with bad data does not hang
227export const pipedDecompressionBadData = {
228 async test() {
229 async function consume(readable) {
230 const readAll = async () => {
231 for await (const _ of readable) {
232 // Should error
233 }
234 };
235 await rejects(readAll(), { message: 'Decompression failed.' });
236 }
237 
238 async function doTest(transform) {
239 const { writable, readable } = new DecompressionStream('gzip');
240 
241 const dest = readable.pipeThrough(transform);
242 
243 const c = consume(dest);
244 
245 const enc = new TextEncoder();
246 const writer = writable.getWriter();
247 await rejects(writer.write(enc.encode('hello world')), {
248 message: 'Decompression failed.',
249 });
250 
251 await c;
252 }
253 
254 await Promise.all([
255 doTest(new IdentityTransformStream()),
256 doTest(new TransformStream()),
257 ]);
258 },
259};
260 
261// Test TransformStream pipeThrough CompressionStream
262export const transformStreamThroughCompression = {
263 async test() {
264 const enc = new TextEncoder();
265 const dec = new TextDecoder();
266 
267 // Create source data
268 const sourceData = 'hello world test data';
269 
270 // Create a TransformStream as source
271 const { readable, writable } = new TransformStream();
272 
273 // Pipe through compression
274 const compressed = readable.pipeThrough(new CompressionStream('gzip'));
275 
276 // Pipe through decompression
277 const decompressed = compressed.pipeThrough(
278 new DecompressionStream('gzip')
279 );
280 
281 // Write source data
282 const writer = writable.getWriter();
283 await writer.write(enc.encode(sourceData));
284 await writer.close();
285 
286 // Read result
287 let result = '';
288 for await (const chunk of decompressed) {
289 result += dec.decode(chunk, { stream: true });
290 }
291 result += dec.decode();
292 
293 strictEqual(result, sourceData);
294 },
295};
296 
297// Test request.body pipeThrough DecompressionStream
298// This test requires a service binding that sends gzip-compressed data
299export const requestBodyPipethrough = {
300 async test(_ctrl, env) {
301 // Fetch compressed data from the service
302 const response = await env.SERVICE.fetch('http://test/compressed');
303 
304 // Pipe the response body through decompression
305 const decompressed = response.body.pipeThrough(
306 new DecompressionStream('gzip')
307 );
308 
309 // Read all decompressed data
310 const dec = new TextDecoder();
311 let result = '';
312 for await (const chunk of decompressed) {
313 result += dec.decode(chunk, { stream: true });
314 }
315 result += dec.decode();
316 
317 strictEqual(result, 'hello world '.repeat(100));
318 },
319};
320 
321// Test transform roundtrip: request body โ†’ DecompressionStream โ†’ TransformStream โ†’ IdentityTransformStream
322export const transformRoundtrip = {
323 async test(_ctrl, env) {
324 // Fetch compressed data from the service
325 const response = await env.SERVICE.fetch('http://test/compressed');
326 
327 const decompression = new DecompressionStream('gzip');
328 const ts = new TransformStream({
329 transform(chunk, controller) {
330 controller.enqueue(chunk);
331 },
332 });
333 const { readable, writable } = new IdentityTransformStream();
334 
335 // Test that piping from a JS-backed TransformStream through an
336 // IdentityTransformStream does not result in a hung pipeTo promise.
337 const pipePromise = response.body
338 .pipeThrough(decompression)
339 .pipeThrough(ts)
340 .pipeTo(writable);
341 
342 // Read the result
343 const dec = new TextDecoder();
344 let result = '';
345 for await (const chunk of readable) {
346 result += dec.decode(chunk, { stream: true });
347 }
348 result += dec.decode();
349 
350 await pipePromise;
351 
352 strictEqual(result, 'hello world '.repeat(100));
353 },
354};
355 
356// Test piping JS-backed stream to internal Response body.
357// Regression test for close signal propagation from JS ReadableStream to internal writable.
358export const compressionPipeline = {
359 async test(_ctrl, env) {
360 const response = await env.SERVICE.fetch('http://test/compressionPipeline');
361 strictEqual(response.status, 200);
362 
363 const decompressed = response.body.pipeThrough(
364 new DecompressionStream('gzip')
365 );
366 
367 const dec = new TextDecoder();
368 let result = '';
369 for await (const chunk of decompressed) {
370 result += dec.decode(chunk, { stream: true });
371 }
372 result += dec.decode();
373 
374 strictEqual(result, 'hello world '.repeat(100));
375 },
376};
377 
378// Test DecompressionStream readable piped to internal Response body.
379export const decompressionPipeline = {
380 async test(_ctrl, env) {
381 const response = await env.SERVICE.fetch(
382 'http://test/decompressionPipeline'
383 );
384 strictEqual(response.status, 200);
385 
386 const dec = new TextDecoder();
387 let result = '';
388 for await (const chunk of response.body) {
389 result += dec.decode(chunk, { stream: true });
390 }
391 result += dec.decode();
392 
393 strictEqual(result, 'hello world '.repeat(100));
394 },
395};
396 
397// Test strictCompression: closing without finishing decompression should error
398export const strictCompressionCloseWithoutFinish = {
399 async test() {
400 const ds = new DecompressionStream('gzip');
401 const writer = ds.writable.getWriter();
402 await rejects(writer.close(), TypeError);
403 },
404};
405 
406// Test strictCompression: trailing data after valid gzip stream should error
407export const strictCompressionTrailingData = {
408 async test() {
409 const ds = new DecompressionStream('gzip');
410 const writer = ds.writable.getWriter();
411 
412 // Gzipped string "FOOBAR", plus a trailing 0xFF byte
413 const trailingStrm = new Uint8Array([
414 0x1f, 0x8b, 0x08, 0x00, 0xf9, 0x05, 0xb7, 0x59, 0x00, 0x03, 0x4b, 0xcb,
415 0xcf, 0x4f, 0x4a, 0x2c, 0x02, 0x00, 0x95, 0x1f, 0xf6, 0x9e, 0x06, 0x00,
416 0x00, 0x00, 0xff,
417 ]);
418 
419 await rejects(writer.write(trailingStrm), TypeError);
420 },
421};
422 
423// =============================================================================
424// Edge case tests for compression streams with chunked data.
425// These tests focus on scenarios where data arrives in multiple chunks,
426// large data handling, and error scenarios.
427//
428// Test inspirations:
429// - Bun: test/js/web/streams/compression.test.ts (round-trip tests, various formats)
430// - Deno: tests/unit/streams_test.ts (compression abort/cancel tests)
431// - Bun: test/js/web/fetch/fetch.stream.test.ts (corrupted data handling)
432// =============================================================================
433 
434// Test compression with many small chunks (1 byte each)
435// Inspired by: Bun test/js/web/streams/compression.test.ts
436export const compressionMultipleSmallChunks = {
437 async test() {
438 const enc = new TextEncoder();
439 const dec = new TextDecoder();
440 const originalText = 'Hello, World! This is a test of chunked compression.';
441 const bytes = enc.encode(originalText);
442 
443 const cs = new CompressionStream('gzip');
444 const writer = cs.writable.getWriter();
445 
446 for (let i = 0; i < bytes.length; i++) {
447 await writer.write(new Uint8Array([bytes[i]]));
448 }
449 await writer.close();
450 
451 const compressedChunks = [];
452 const reader = cs.readable.getReader();
453 while (true) {
454 const { value, done } = await reader.read();
455 if (done) break;
456 compressedChunks.push(value);
457 }
458 
459 const ds = new DecompressionStream('gzip');
460 const dsWriter = ds.writable.getWriter();
461 for (const chunk of compressedChunks) {
462 await dsWriter.write(chunk);
463 }
464 await dsWriter.close();
465 
466 let result = '';
467 const dsReader = ds.readable.getReader();
468 while (true) {
469 const { value, done } = await dsReader.read();
470 if (done) break;
471 result += dec.decode(value, { stream: true });
472 }
473 result += dec.decode();
474 
475 strictEqual(result, originalText);
476 },
477};
478 
479// Test decompression with chunked input (compressed data arrives in pieces)
480// Inspired by: Bun test/js/web/fetch/fetch.stream.test.ts
481export const decompressionChunkedInput = {
482 async test() {
483 const enc = new TextEncoder();
484 const dec = new TextDecoder();
485 const originalText = 'Test data for chunked decompression testing.';
486 
487 const cs = new CompressionStream('deflate');
488 const csWriter = cs.writable.getWriter();
489 await csWriter.write(enc.encode(originalText));
490 await csWriter.close();
491 
492 const compressedChunks = [];
493 const csReader = cs.readable.getReader();
494 while (true) {
495 const { value, done } = await csReader.read();
496 if (done) break;
497 compressedChunks.push(value);
498 }
499 
500 const totalLength = compressedChunks.reduce(
501 (sum, chunk) => sum + chunk.length,
502 0
503 );
504 const compressed = new Uint8Array(totalLength);
505 let offset = 0;
506 for (const chunk of compressedChunks) {
507 compressed.set(chunk, offset);
508 offset += chunk.length;
509 }
510 
511 const ds = new DecompressionStream('deflate');
512 const dsWriter = ds.writable.getWriter();
513 
514 for (let i = 0; i < compressed.length; i += 2) {
515 const end = Math.min(i + 2, compressed.length);
516 await dsWriter.write(compressed.slice(i, end));
517 }
518 await dsWriter.close();
519 
520 let result = '';
521 const dsReader = ds.readable.getReader();
522 while (true) {
523 const { value, done } = await dsReader.read();
524 if (done) break;
525 result += dec.decode(value, { stream: true });
526 }
527 result += dec.decode();
528 
529 strictEqual(result, originalText);
530 },
531};
532 
533// Test compression with large data (100KB+)
534// Inspired by: Bun bench/snippets/compression-streams.mjs
535export const compressionLargeData = {
536 async test() {
537 const size = 100 * 1024;
538 const data = new Uint8Array(size);
539 for (let i = 0; i < size; i++) {
540 data[i] = i % 256;
541 }
542 
543 const cs = new CompressionStream('gzip');
544 const csWriter = cs.writable.getWriter();
545 await csWriter.write(data);
546 await csWriter.close();
547 
548 const compressedChunks = [];
549 const csReader = cs.readable.getReader();
550 while (true) {
551 const { value, done } = await csReader.read();
552 if (done) break;
553 compressedChunks.push(value);
554 }
555 
556 const compressedLength = compressedChunks.reduce(
557 (sum, chunk) => sum + chunk.length,
558 0
559 );
560 ok(compressedLength < size, 'Compressed should be smaller than original');
561 
562 const ds = new DecompressionStream('gzip');
563 const dsWriter = ds.writable.getWriter();
564 for (const chunk of compressedChunks) {
565 await dsWriter.write(chunk);
566 }
567 await dsWriter.close();
568 
569 const decompressedChunks = [];
570 const dsReader = ds.readable.getReader();
571 while (true) {
572 const { value, done } = await dsReader.read();
573 if (done) break;
574 decompressedChunks.push(value);
575 }
576 
577 const decompressedLength = decompressedChunks.reduce(
578 (sum, chunk) => sum + chunk.length,
579 0
580 );
581 strictEqual(decompressedLength, size);
582 
583 const decompressed = new Uint8Array(decompressedLength);
584 let dOffset = 0;
585 for (const chunk of decompressedChunks) {
586 decompressed.set(chunk, dOffset);
587 dOffset += chunk.length;
588 }
589 
590 for (let i = 0; i < size; i++) {
591 strictEqual(decompressed[i], i % 256, `Byte at ${i} should match`);
592 }
593 },
594};
595 
596// Test decompression with completely invalid data should error
597// Inspired by: Bun test/js/web/fetch/fetch.stream.test.ts (corrupted data handling)
598export const decompressionTruncated = {
599 async test() {
600 const invalidData = new Uint8Array([0x00, 0x01, 0x02, 0x03, 0x04, 0x05]);
601 
602 const ds = new DecompressionStream('gzip');
603 const dsWriter = ds.writable.getWriter();
604 const dsReader = ds.readable.getReader();
605 
606 await rejects(async () => {
607 await dsWriter.write(invalidData);
608 await dsWriter.close();
609 
610 while (true) {
611 const { done } = await dsReader.read();
612 if (done) break;
613 }
614 });
615 },
616};
617 
618// Test empty stream through compression
619// Inspired by: Bun test/js/web/streams/compression.test.ts
620export const compressionEmptyStream = {
621 async test() {
622 const cs = new CompressionStream('deflate');
623 const writer = cs.writable.getWriter();
624 await writer.close();
625 
626 const chunks = [];
627 const reader = cs.readable.getReader();
628 while (true) {
629 const { value, done } = await reader.read();
630 if (done) break;
631 chunks.push(value);
632 }
633 
634 const ds = new DecompressionStream('deflate');
635 const dsWriter = ds.writable.getWriter();
636 for (const chunk of chunks) {
637 await dsWriter.write(chunk);
638 }
639 await dsWriter.close();
640 
641 const dsReader = ds.readable.getReader();
642 const result = await dsReader.read();
643 ok(result.done);
644 },
645};
646 
647// Test empty decompression
648// Inspired by: Bun test/js/web/streams/compression.test.ts
649export const decompressionEmptyStream = {
650 async test() {
651 // Create valid empty compressed data first
652 const cs = new CompressionStream('gzip');
653 const csWriter = cs.writable.getWriter();
654 await csWriter.close();
655 
656 const compressedChunks = [];
657 const csReader = cs.readable.getReader();
658 while (true) {
659 const { value, done } = await csReader.read();
660 if (done) break;
661 compressedChunks.push(value);
662 }
663 
664 // Now decompress
665 const ds = new DecompressionStream('gzip');
666 const dsWriter = ds.writable.getWriter();
667 for (const chunk of compressedChunks) {
668 await dsWriter.write(chunk);
669 }
670 await dsWriter.close();
671 
672 // Verify empty result
673 const dsReader = ds.readable.getReader();
674 const result = await dsReader.read();
675 ok(result.done);
676 },
677};
678 
679// Test CompressionStream piped to IdentityTransformStream
680// Inspired by: workerd compression-streams-test.js
681export const compressionPipeToIdentity = {
682 async test() {
683 const enc = new TextEncoder();
684 const dec = new TextDecoder();
685 const originalText = 'Piping compression to identity transform.';
686 
687 const source = new ReadableStream({
688 start(controller) {
689 controller.enqueue(enc.encode(originalText));
690 controller.close();
691 },
692 });
693 
694 const compressed = source
695 .pipeThrough(new CompressionStream('deflate'))
696 .pipeThrough(new IdentityTransformStream());
697 
698 const decompressed = compressed.pipeThrough(
699 new DecompressionStream('deflate')
700 );
701 
702 let result = '';
703 for await (const chunk of decompressed) {
704 result += dec.decode(chunk, { stream: true });
705 }
706 result += dec.decode();
707 
708 strictEqual(result, originalText);
709 },
710};
711 
712// Test all supported formats work with chunked data
713// Inspired by: Bun test/js/web/streams/compression.test.ts
714export const decompressionAllFormats = {
715 async test() {
716 const enc = new TextEncoder();
717 const dec = new TextDecoder();
718 const formats = ['gzip', 'deflate', 'deflate-raw'];
719 const originalText = 'Testing all compression formats with chunked data!';
720 
721 for (const format of formats) {
722 // Compress
723 const cs = new CompressionStream(format);
724 const csWriter = cs.writable.getWriter();
725 
726 // Write in chunks
727 const bytes = enc.encode(originalText);
728 for (let i = 0; i < bytes.length; i += 5) {
729 await csWriter.write(bytes.slice(i, Math.min(i + 5, bytes.length)));
730 }
731 await csWriter.close();
732 
733 // Collect compressed
734 const compressedChunks = [];
735 const csReader = cs.readable.getReader();
736 while (true) {
737 const { value, done } = await csReader.read();
738 if (done) break;
739 compressedChunks.push(value);
740 }
741 
742 // Decompress
743 const ds = new DecompressionStream(format);
744 const dsWriter = ds.writable.getWriter();
745 for (const chunk of compressedChunks) {
746 await dsWriter.write(chunk);
747 }
748 await dsWriter.close();
749 
750 // Verify
751 let result = '';
752 const dsReader = ds.readable.getReader();
753 while (true) {
754 const { value, done } = await dsReader.read();
755 if (done) break;
756 result += dec.decode(value, { stream: true });
757 }
758 result += dec.decode();
759 
760 strictEqual(
761 result,
762 originalText,
763 `Format ${format} should round-trip correctly`
764 );
765 }
766 },
767};
768 
769export default {
770 async fetch(request, env) {
771 if (request.url.includes('/compressed')) {
772 const data = 'hello world '.repeat(100);
773 const enc = new TextEncoder();
774 const { readable, writable } = new CompressionStream('gzip');
775 const writer = writable.getWriter();
776 await writer.write(enc.encode(data));
777 await writer.close();
778 return new Response(readable, {
779 headers: { 'Content-Encoding': 'gzip' },
780 });
781 }
782 if (request.url.includes('/stream')) {
783 return new Response('hello world '.repeat(100));
784 }
785 if (request.url.includes('/compressionPipeline')) {
786 // Pipe through TransformStream โ†’ CompressionStream โ†’ Response body (internal writable)
787 const response = await env.SERVICE.fetch('http://test/stream');
788 const { readable, writable } = new TransformStream();
789 response.body.pipeTo(writable);
790 const compressed = readable.pipeThrough(new CompressionStream('gzip'));
791 return new Response(compressed, { encodeBody: 'manual' });
792 }
793 if (request.url.includes('/decompressionPipeline')) {
794 // Pipe DecompressionStream readable โ†’ Response body (internal writable)
795 const response = await env.SERVICE.fetch('http://test/compressed');
796 const decompressed = response.body.pipeThrough(
797 new DecompressionStream('gzip')
798 );
799 return new Response(decompressed);
800 }
801 return new Response('Not found', { status: 404 });
802 },
803};