Skip to content
Open
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
25 changes: 19 additions & 6 deletions src/utils/create-stream-promise.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,16 @@ const createStreamPromise = function (
let dataLength = 0

let timeoutId: NodeJS.Timeout | null = null
const clearTimer = () => {
if (timeoutId) {
clearTimeout(timeoutId)
timeoutId = null
}
}
if (timeoutSeconds != null && Number.isFinite(timeoutSeconds)) {
timeoutId = setTimeout(() => {
data = null
clearTimer()
reject(new Error('Request timed out waiting for body'))
}, timeoutSeconds * SEC_TO_MILLISEC)
}
Expand All @@ -31,6 +38,7 @@ const createStreamPromise = function (
dataLength += chunk.length
if (dataLength > bytesLimit) {
data = null
clearTimer()
reject(new Error('Stream body too big'))
} else {
data.push(chunk)
Expand All @@ -40,16 +48,21 @@ const createStreamPromise = function (
stream.on('error', function onError(error) {
data = null
reject(error)
if (timeoutId) {
clearTimeout(timeoutId)
clearTimer()
})
stream.on('close', () => {
if (data) {
data = null
clearTimer()
reject(new Error('Stream closed before body completed'))
}
})
stream.on('end', function onEnd() {
if (timeoutId) {
clearTimeout(timeoutId)
}
clearTimer()
if (data) {
resolve(Buffer.concat(data))
const body = Buffer.concat(data)
data = null
resolve(body)
}
})
})
Expand Down
84 changes: 84 additions & 0 deletions tests/unit/utils/create-stream-promise.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
import { PassThrough } from 'node:stream'
import { afterEach, describe, expect, test, vi } from 'vitest'
import createStreamPromise from '../../../src/utils/create-stream-promise.js'

afterEach(() => vi.useRealTimers())

describe('createStreamPromise', () => {
test('clears the request timeout when the size limit rejects a live stream', async () => {
vi.useFakeTimers()
const stream = new PassThrough()
const body = createStreamPromise(stream, 30, 2)
stream.write(Buffer.from('abc'))
await expect(body).rejects.toThrow('Stream body too big')
expect(vi.getTimerCount()).toBe(0)
})

test('rejects a closed stream promptly instead of waiting for the request timeout', async () => {
vi.useFakeTimers()
const stream = new PassThrough()
const body = createStreamPromise(stream, 30)
const rejected = body.catch((error: unknown) => error)
stream.emit('close')
await vi.advanceTimersByTimeAsync(0)
expect(vi.getTimerCount()).toBe(0)
expect(await rejected).toBeInstanceOf(Error)
expect(((await rejected) as Error).message).toMatch(
/Stream closed before body completed|Request timed out waiting for body/,
)
})

test('rejects premature closure even without a timeout', async () => {
const stream = new PassThrough()
const body = createStreamPromise(stream, Number.POSITIVE_INFINITY)
const rejected = body.catch((error: unknown) => error)
stream.destroy()
expect(await rejected).toBeInstanceOf(Error)
expect(((await rejected) as Error).message).toMatch(
/Stream closed before body completed|Request timed out waiting for body/,
)
})

test('ignores later events after a size rejection', async () => {
vi.useFakeTimers()
const stream = new PassThrough()
const body = createStreamPromise(stream, 30, 2)
stream.write(Buffer.from('abc'))
await expect(body).rejects.toThrow('Stream body too big')
stream.emit('error', new Error('later error'))
stream.emit('close')
stream.end()
expect(vi.getTimerCount()).toBe(0)
})

test('joins chunks when the stream ends normally', async () => {
const stream = new PassThrough()
const body = createStreamPromise(stream, 30)
stream.write(Buffer.from('ab'))
stream.end(Buffer.from('cd'))
await expect(body).resolves.toEqual(Buffer.from('abcd'))
})

test('preserves stream errors and clears the timeout', async () => {
vi.useFakeTimers()
const stream = new PassThrough()
const body = createStreamPromise(stream, 30)
const error = new Error('read failed')
stream.emit('error', error)
await expect(body).rejects.toBe(error)
expect(vi.getTimerCount()).toBe(0)
})

test('still rejects a stalled stream at the timeout', async () => {
vi.useFakeTimers()
const stream = new PassThrough()
const body = createStreamPromise(stream, 1)
const rejected = body.catch((error: unknown) => error)
await vi.advanceTimersByTimeAsync(1000)
expect(await rejected).toBeInstanceOf(Error)
expect(((await rejected) as Error).message).toMatch(
/Stream closed before body completed|Request timed out waiting for body/,
)
expect(vi.getTimerCount()).toBe(0)
})
})