diff --git a/src/interceptors/fetch/utils/brotli-decompress.ts b/src/interceptors/fetch/utils/brotli-decompress.ts index dbfa622e..8ab3d521 100644 --- a/src/interceptors/fetch/utils/brotli-decompress.ts +++ b/src/interceptors/fetch/utils/brotli-decompress.ts @@ -1,31 +1,83 @@ import zlib from 'node:zlib' +import type { Transformer } from 'node:stream/web' -export class BrotliDecompressionStream extends TransformStream { +export class BrotliDecompressionStream extends TransformStream< + Uint8Array, + Uint8Array +> { constructor() { const decompress = zlib.createBrotliDecompress({ flush: zlib.constants.BROTLI_OPERATION_FLUSH, finishFlush: zlib.constants.BROTLI_OPERATION_FLUSH, }) - super({ - async transform(chunk, controller) { - const buffer = Buffer.from(chunk) + /** + * @note Typed via Node.js as the DOM `Transformer` type + * does not declare the `cancel` callback yet. + */ + const transformer: Transformer = { + start(controller) { + /** + * @note Forward every decompressed chunk to the stream. + * A single input chunk can produce multiple output chunks + * (Node.js emits Brotli output in 16 KiB chunks by default). + * @see https://github.com/mswjs/interceptors/issues/798 + */ + decompress.on('data', (chunk: Buffer) => { + // Ignore any in-flight output once the consumer cancelled the stream. + if (decompress.destroyed) { + return + } - const decompressed = await new Promise((resolve, reject) => { - decompress.write(buffer, (error) => { - if (error) reject(error) - }) - - decompress.flush() - decompress.once('data', (data) => resolve(data)) - decompress.once('error', (error) => reject(error)) - decompress.once('end', () => controller.terminate()) - }).catch((error) => { + try { + controller.enqueue(new Uint8Array(chunk)) + } catch { + /** + * @note The readable side has been cancelled. If that happens + * after `flush()` has started, the `cancel()` callback is never + * invoked (per the Streams spec), so stop decompression here. + */ + decompress.destroy() + } + }) + decompress.once('error', (error) => { controller.error(error) }) + }, + transform(chunk) { + const writePromise = Promise.withResolvers() + + decompress.write(chunk, (error) => { + if (error) { + writePromise.reject(error) + } else { + writePromise.resolve() + } + }) - controller.enqueue(decompressed) + return writePromise.promise }, - }) + flush() { + const flushPromise = Promise.withResolvers() + + decompress.once('end', flushPromise.resolve) + // Settle the flush if decompression is stopped by a cancellation. + decompress.once('close', flushPromise.resolve) + decompress.once('error', flushPromise.reject) + decompress.end() + + return flushPromise.promise + }, + cancel() { + /** + * @note Release the zlib handle once the consumer cancels + * the readable side. Nothing can be enqueued into a cancelled + * stream, so any pending decompression output must be dropped. + */ + decompress.destroy() + }, + } + + super(transformer) } } diff --git a/src/interceptors/fetch/utils/decompression.test.ts b/src/interceptors/fetch/utils/decompression.test.ts new file mode 100644 index 00000000..5daf235c --- /dev/null +++ b/src/interceptors/fetch/utils/decompression.test.ts @@ -0,0 +1,146 @@ +// @vitest-environment node +import { brotliCompressSync } from 'node:zlib' +import { setTimeout } from 'node:timers/promises' +import { decompressResponse, isCompressedResponse } from './decompression' +import { BrotliDecompressionStream } from './brotli-decompress' + +describe('isCompressedResponse', () => { + it('returns false for a response without a body', () => { + const response = new Response(null, { + headers: { 'content-encoding': 'gzip' }, + }) + + expect(isCompressedResponse(response)).toBe(false) + }) + + it('returns false for a response without "content-encoding"', () => { + expect(isCompressedResponse(new Response('hello world'))).toBe(false) + }) + + it('returns false for a response with an empty "content-encoding"', () => { + const response = new Response('hello world', { + headers: { 'content-encoding': '' }, + }) + + expect(isCompressedResponse(response)).toBe(false) + }) + + it('returns true for a response with a body and "content-encoding"', () => { + const response = new Response('hello world', { + headers: { 'content-encoding': 'gzip' }, + }) + + expect(isCompressedResponse(response)).toBe(true) + }) +}) + +describe('decompressResponse', () => { + it('throws when given a non-compressed response', () => { + expect(() => { + decompressResponse(new Response('hello world')) + }).toThrow(/Failed to decompress a response/) + }) + + /** + * @see https://github.com/mswjs/interceptors/issues/798 + * @note Node.js emits decompressed Brotli output in 16 KiB chunks. + * A body larger than that must be decompressed in full, not truncated. + */ + it('decompresses a "br" body larger than a single zlib chunk', async () => { + const expectedBody = Array.from( + { length: 10_000 }, + (_, index) => `line ${index}: ${crypto.randomUUID()}` + ).join('\n') + const response = new Response(brotliCompressSync(expectedBody), { + headers: { 'content-encoding': 'br' }, + }) + + const actualBody = await new Response( + decompressResponse(response) + ).text() + + expect(actualBody.length).toBe(expectedBody.length) + expect(actualBody).toBe(expectedBody) + }) + + it('decompresses a "br" body streamed in multiple chunks', async () => { + const expectedBody = Array.from( + { length: 10_000 }, + (_, index) => `line ${index}: ${crypto.randomUUID()}` + ).join('\n') + const compressedBody = brotliCompressSync(expectedBody) + const chunkSize = 1024 + const compressedStream = new ReadableStream({ + start(controller) { + for ( + let offset = 0; + offset < compressedBody.byteLength; + offset += chunkSize + ) { + controller.enqueue( + new Uint8Array(compressedBody.subarray(offset, offset + chunkSize)) + ) + } + controller.close() + }, + }) + const response = new Response(compressedStream, { + headers: { 'content-encoding': 'br' }, + }) + + const actualBody = await new Response( + decompressResponse(response) + ).text() + + expect(actualBody).toBe(expectedBody) + }) + + it('errors when given a corrupted "br" body', async () => { + const response = new Response(new Uint8Array(16).fill(0xff), { + headers: { 'content-encoding': 'br' }, + }) + + await expect( + new Response(decompressResponse(response)).text() + ).rejects.toThrow() + }) +}) + +describe('BrotliDecompressionStream', () => { + it('drops pending output once the readable side is cancelled', async () => { + const compressedBody = brotliCompressSync(Buffer.alloc(200_000, 'a')) + const stream = new BrotliDecompressionStream() + const writer = stream.writable.getWriter() + const reader = stream.readable.getReader() + + writer.write(new Uint8Array(compressedBody)).catch(() => {}) + // Read the first decompressed chunk while zlib still has output pending. + await expect(reader.read()).resolves.toHaveProperty('done', false) + + await reader.cancel('consumer cancelled') + + // Pending zlib output must not be enqueued into the cancelled stream. + // That would throw an uncaught "Invalid state" error, failing this test. + await setTimeout(50) + }) + + it('settles a cancellation that happens during flush', async () => { + const compressedBody = brotliCompressSync(Buffer.alloc(200_000, 'a')) + const stream = new BrotliDecompressionStream() + const writer = stream.writable.getWriter() + const reader = stream.readable.getReader() + + writer.write(new Uint8Array(compressedBody)).catch(() => {}) + await expect(reader.read()).resolves.toHaveProperty('done', false) + // Closing the writable side starts the flush. + const closePromise = writer.close() + + // Cancelling during flush does not invoke the "cancel" callback. + // The stream must still settle and must not enqueue pending output. + await expect(reader.cancel('consumer cancelled')).resolves.toBeUndefined() + // Per the Streams spec, the writable side errors with the cancel reason. + await expect(closePromise).rejects.toBe('consumer cancelled') + await setTimeout(50) + }) +}) +