Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
84 changes: 68 additions & 16 deletions src/interceptors/fetch/utils/brotli-decompress.ts
Original file line number Diff line number Diff line change
@@ -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<Uint8Array, Uint8Array> = {
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<Buffer>((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()
}
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.
decompress.once('error', (error) => {
controller.error(error)
})
},
transform(chunk) {
const writePromise = Promise.withResolvers<void>()

decompress.write(chunk, (error) => {
if (error) {
writePromise.reject(error)
} else {
writePromise.resolve()
}
})

controller.enqueue(decompressed)
return writePromise.promise
},
})
flush() {
const flushPromise = Promise.withResolvers<void>()

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()
Comment thread
kettanaito marked this conversation as resolved.
},
}

super(transformer)
}
}
146 changes: 146 additions & 0 deletions src/interceptors/fetch/utils/decompression.test.ts
Original file line number Diff line number Diff line change
@@ -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<Uint8Array>({
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)
})
})

Loading