From 115520bb0499f83e77b5891251a6c8c40b9c7f9b Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Mon, 28 Sep 2026 14:40:08 +0200 Subject: [PATCH 1/3] fix(fetch): prevent truncating compressed responses to the first chunk --- .../fetch/utils/brotli-decompress.ts | 51 ++++++--- .../fetch/utils/decompression.test.ts | 105 ++++++++++++++++++ 2 files changed, 141 insertions(+), 15 deletions(-) create mode 100644 src/interceptors/fetch/utils/decompression.test.ts diff --git a/src/interceptors/fetch/utils/brotli-decompress.ts b/src/interceptors/fetch/utils/brotli-decompress.ts index dbfa622e5..1736f8894 100644 --- a/src/interceptors/fetch/utils/brotli-decompress.ts +++ b/src/interceptors/fetch/utils/brotli-decompress.ts @@ -1,6 +1,9 @@ import zlib from 'node:zlib' -export class BrotliDecompressionStream extends TransformStream { +export class BrotliDecompressionStream extends TransformStream< + Uint8Array, + Uint8Array +> { constructor() { const decompress = zlib.createBrotliDecompress({ flush: zlib.constants.BROTLI_OPERATION_FLUSH, @@ -8,23 +11,41 @@ export class BrotliDecompressionStream extends TransformStream { }) super({ - async transform(chunk, controller) { - const buffer = Buffer.from(chunk) - - 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) => { + 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) => { + controller.enqueue(new Uint8Array(chunk)) + }) + 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() + } + }) + + return writePromise.promise + }, + flush() { + const flushPromise = Promise.withResolvers() + + decompress.once('end', flushPromise.resolve) + decompress.once('error', flushPromise.reject) + decompress.end() - controller.enqueue(decompressed) + return flushPromise.promise }, }) } diff --git a/src/interceptors/fetch/utils/decompression.test.ts b/src/interceptors/fetch/utils/decompression.test.ts new file mode 100644 index 000000000..9372bdeae --- /dev/null +++ b/src/interceptors/fetch/utils/decompression.test.ts @@ -0,0 +1,105 @@ +// @vitest-environment node +import { brotliCompressSync } from 'node:zlib' +import { decompressResponse, isCompressedResponse } from './decompression' + +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() + }) +}) From 3c0fb5b133d3360047319bf26ef2922e6a658539 Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Mon, 28 Sep 2026 14:49:48 +0200 Subject: [PATCH 2/3] fix: handle brotli stream cancellation --- .../fetch/utils/brotli-decompress.ts | 24 +++++++++++++++++-- .../fetch/utils/decompression.test.ts | 22 +++++++++++++++++ 2 files changed, 44 insertions(+), 2 deletions(-) diff --git a/src/interceptors/fetch/utils/brotli-decompress.ts b/src/interceptors/fetch/utils/brotli-decompress.ts index 1736f8894..106d471c3 100644 --- a/src/interceptors/fetch/utils/brotli-decompress.ts +++ b/src/interceptors/fetch/utils/brotli-decompress.ts @@ -1,4 +1,5 @@ import zlib from 'node:zlib' +import type { Transformer } from 'node:stream/web' export class BrotliDecompressionStream extends TransformStream< Uint8Array, @@ -10,7 +11,11 @@ export class BrotliDecompressionStream extends TransformStream< finishFlush: zlib.constants.BROTLI_OPERATION_FLUSH, }) - super({ + /** + * @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. @@ -19,6 +24,11 @@ export class BrotliDecompressionStream extends TransformStream< * @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 + } + controller.enqueue(new Uint8Array(chunk)) }) decompress.once('error', (error) => { @@ -47,6 +57,16 @@ export class BrotliDecompressionStream extends TransformStream< 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 index 9372bdeae..ee0322096 100644 --- a/src/interceptors/fetch/utils/decompression.test.ts +++ b/src/interceptors/fetch/utils/decompression.test.ts @@ -1,6 +1,8 @@ // @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', () => { @@ -103,3 +105,23 @@ describe('decompressResponse', () => { ).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) + }) +}) + From 9c41a44bb5e9a25dc6e302cdf0824923df9e0c3d Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Mon, 28 Sep 2026 15:00:48 +0200 Subject: [PATCH 3/3] fix: settle a cancellation that happens during flush --- .../fetch/utils/brotli-decompress.ts | 13 ++++++++++++- .../fetch/utils/decompression.test.ts | 19 +++++++++++++++++++ 2 files changed, 31 insertions(+), 1 deletion(-) diff --git a/src/interceptors/fetch/utils/brotli-decompress.ts b/src/interceptors/fetch/utils/brotli-decompress.ts index 106d471c3..8ab3d521c 100644 --- a/src/interceptors/fetch/utils/brotli-decompress.ts +++ b/src/interceptors/fetch/utils/brotli-decompress.ts @@ -29,7 +29,16 @@ export class BrotliDecompressionStream extends TransformStream< return } - controller.enqueue(new Uint8Array(chunk)) + 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) @@ -52,6 +61,8 @@ export class BrotliDecompressionStream extends TransformStream< 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() diff --git a/src/interceptors/fetch/utils/decompression.test.ts b/src/interceptors/fetch/utils/decompression.test.ts index ee0322096..5daf235ce 100644 --- a/src/interceptors/fetch/utils/decompression.test.ts +++ b/src/interceptors/fetch/utils/decompression.test.ts @@ -123,5 +123,24 @@ describe('BrotliDecompressionStream', () => { // 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) + }) })