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
6 changes: 5 additions & 1 deletion src/listener.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,11 @@ const drainIncoming = (incoming: IncomingMessage | Http2ServerRequest): void =>
cleanup()
const socket = incoming.socket
if (socket && !socket.destroyed) {
socket.destroySoon()
if (typeof socket.destroySoon === 'function') {
socket.destroySoon()
} else if (typeof socket.destroy === 'function') {
socket.destroy()
}
}
}

Expand Down
84 changes: 84 additions & 0 deletions test/listener.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
import { EventEmitter } from 'node:events'
import { createServer } from 'node:http'
import type { IncomingMessage, ServerResponse } from 'node:http'
import { Readable } from 'node:stream'
import { getRequestListener } from '../src/listener'
import { GlobalRequest, Request as LightweightRequest, RequestError } from '../src/request'
import { GlobalResponse, Response as LightweightResponse } from '../src/response'
Expand Down Expand Up @@ -542,6 +545,87 @@ describe('Abort request - cacheable response path', () => {
})
})

describe('Non-standard incoming request', () => {
class MockSocket extends EventEmitter {
remoteAddress = '127.0.0.1'
remotePort = 44936
}

class MockIncomingMessage extends Readable {
method = 'POST'
url = '/'
headers = { host: 'localhost' }
rawHeaders = ['host', 'localhost']

constructor(readonly socket: EventEmitter = new MockSocket()) {
super()
}

_read() {
// The body is never pushed and never ends, so draining cannot complete
// and the drain timeout is guaranteed to fire.
}
}

class MockServerResponse extends EventEmitter {
headersSent = false
writableFinished = false

writeHead() {
this.headersSent = true
return this
}

end() {
this.writableFinished = true
this.emit('finish')
this.emit('close')
return this
}
}

it('Should not throw when the drain timeout fires on a socket without destroySoon', async () => {
vi.useFakeTimers()

try {
const requestListener = getRequestListener(() => new LightweightResponse('ok'))
await requestListener(
new MockIncomingMessage() as unknown as IncomingMessage,
new MockServerResponse() as unknown as ServerResponse
)

// The drain timeout fires on a timer, so a throw here is an uncatchable
// uncaughtException for the caller.
expect(() => vi.advanceTimersByTime(1_000)).not.toThrow()
} finally {
vi.useRealTimers()
}
})

it('Should destroy a socket that implements destroy but not destroySoon', async () => {
vi.useFakeTimers()

try {
// Duplex-based fakes provide the standard stream teardown method only.
class DuplexLikeSocket extends MockSocket {
destroy = vi.fn()
}
const socket = new DuplexLikeSocket()

const requestListener = getRequestListener(() => new LightweightResponse('ok'))
await requestListener(
new MockIncomingMessage(socket) as unknown as IncomingMessage,
new MockServerResponse() as unknown as ServerResponse
)
vi.advanceTimersByTime(1_000)

expect(socket.destroy).toHaveBeenCalledTimes(1)
} finally {
vi.useRealTimers()
}
})
})

describe('overrideGlobalObjects', () => {
const fetchCallback = vi.fn()

Expand Down