Update license

This commit is contained in:
HarithaVattikuti
2026-10-08 11:32:19 -05:00
parent 3e7e9e4543
commit 0d7d4683c1
3 changed files with 907 additions and 375 deletions

View File

@@ -1,6 +1,6 @@
--- ---
name: undici name: undici
version: 6.28.0 version: 6.29.0
type: npm type: npm
summary: An HTTP/1.1 client, written from scratch for Node.js summary: An HTTP/1.1 client, written from scratch for Node.js
homepage: https://undici.nodejs.org homepage: https://undici.nodejs.org

View File

@@ -15982,11 +15982,77 @@ class Request {
} }
} }
onUpgrade (statusCode, headers, socket) { /**
* @param {number|null} statusCode
* @param {Buffer[]|null} headers
* @param {import('node:stream').Duplex} socket
* @param {string} [statusText]
*/
onUpgrade (statusCode, headers, socket, statusText = '') {
this.onFinally()
assert(!this.aborted) assert(!this.aborted)
assert(!this.completed) assert(!this.completed)
return this[kHandler].onUpgrade(statusCode, headers, socket) if (statusCode !== null) {
this.#publishUpgradeHeaders(statusCode, headers, statusText)
}
const result = this[kHandler].onUpgrade(statusCode, headers, socket)
if (!this.aborted) {
this.completed = true
if (statusCode !== null) {
this.#publishUpgradeTrailers()
}
}
return result
}
/**
* @param {number} statusCode
* @param {import('node:http2').IncomingHttpHeaders} headers
* @param {(headers: import('node:http2').IncomingHttpHeaders) => Buffer[]} parseHeaders
* @param {string} [statusText]
*/
onUpgradeResponse (statusCode, headers, parseHeaders, statusText = '') {
assert(!this.aborted)
assert(this.completed)
if (channels.headers.hasSubscribers) {
this.#publishUpgradeHeaders(statusCode, parseHeaders(headers), statusText)
}
this.#publishUpgradeTrailers()
}
/**
* @param {Error} error
*/
onUpgradeError (error) {
assert(!this.aborted)
assert(this.completed)
if (channels.error.hasSubscribers) {
channels.error.publish({ request: this, error })
}
}
/**
* @param {number} statusCode
* @param {Buffer[]} headers
* @param {string} statusText
*/
#publishUpgradeHeaders (statusCode, headers, statusText) {
if (channels.headers.hasSubscribers) {
channels.headers.publish({ request: this, response: { statusCode, headers, statusText } })
}
}
#publishUpgradeTrailers () {
if (channels.trailers.hasSubscribers) {
channels.trailers.publish({ request: this, trailers: [] })
}
} }
onComplete (trailers) { onComplete (trailers) {
@@ -17885,7 +17951,7 @@ class Parser {
} }
onUpgrade (head) { onUpgrade (head) {
const { upgrade, client, socket, headers, statusCode } = this const { upgrade, client, socket, headers, statusCode, statusText } = this
assert(upgrade) assert(upgrade)
assert(client[kSocket] === socket) assert(client[kSocket] === socket)
@@ -17920,9 +17986,10 @@ class Parser {
client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade')) client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade'))
try { try {
request.onUpgrade(statusCode, headers, socket) request.onUpgrade(statusCode, headers, socket, statusText)
} catch (err) { } catch (error) {
util.destroy(socket, err) util.errorRequest(client, request, error)
util.destroy(socket, error)
} }
client[kResume]() client[kResume]()
@@ -18329,7 +18396,7 @@ async function connectH1 (client, socket) {
function clearIdleSocketValidation (socket) { function clearIdleSocketValidation (socket) {
if (socket[kIdleSocketValidationTimeout]) { if (socket[kIdleSocketValidationTimeout]) {
clearTimeout(socket[kIdleSocketValidationTimeout]) clearImmediate(socket[kIdleSocketValidationTimeout])
socket[kIdleSocketValidationTimeout] = null socket[kIdleSocketValidationTimeout] = null
} }
@@ -18338,15 +18405,23 @@ function clearIdleSocketValidation (socket) {
function scheduleIdleSocketValidation (client, socket) { function scheduleIdleSocketValidation (client, socket) {
socket[kIdleSocketValidation] = 1 socket[kIdleSocketValidation] = 1
socket[kIdleSocketValidationTimeout] = setTimeout(() => { // Yield to the check phase (after poll) so unsolicited bytes / FIN / RST
// already pending on this idle keep-alive socket are processed before the
// next request is written (GHSA-35p6-xmwp-9g52).
//
// setTimeout(0) pays Node's ~1ms timer floor on every sequential reuse
// (#5493). setImmediate avoids that, but an *unref'd* Immediate lets poll
// block for ~500ms when the event loop is otherwise idle (#5600 / #5606).
// A ref'd Immediate both keeps the pending request alive and makes poll
// return immediately — the hybrid those issues asked for.
socket[kIdleSocketValidationTimeout] = setImmediate(() => {
socket[kIdleSocketValidationTimeout] = null socket[kIdleSocketValidationTimeout] = null
socket[kIdleSocketValidation] = 2 socket[kIdleSocketValidation] = 2
if (client[kSocket] === socket && !socket.destroyed) { if (client[kSocket] === socket && !socket.destroyed) {
client[kResume]() client[kResume]()
} }
}, 0) })
socket[kIdleSocketValidationTimeout].unref?.()
} }
/** /**
@@ -18495,12 +18570,22 @@ function writeH1 (client, request) {
const socket = client[kSocket] const socket = client[kSocket]
clearIdleSocketValidation(socket) clearIdleSocketValidation(socket)
const abort = (err) => { /**
if (request.aborted || request.completed) { * @param {Error} [error]
*/
const abort = (error) => {
if (request.aborted) {
return return
} }
util.errorRequest(client, request, err || new RequestAbortedError()) if (request.completed) {
if (request.upgrade || request.method === 'CONNECT') {
util.destroy(socket, new InformationalError('aborted'))
}
return
}
util.errorRequest(client, request, error || new RequestAbortedError())
util.destroy(body) util.destroy(body)
util.destroy(socket, new InformationalError('aborted')) util.destroy(socket, new InformationalError('aborted'))
@@ -18957,6 +19042,7 @@ module.exports = connectH1
const assert = __nccwpck_require__(4589) const assert = __nccwpck_require__(4589)
const { errorMonitor } = __nccwpck_require__(8474)
const { pipeline } = __nccwpck_require__(7075) const { pipeline } = __nccwpck_require__(7075)
const util = __nccwpck_require__(3440) const util = __nccwpck_require__(3440)
const { const {
@@ -19033,6 +19119,15 @@ function parseH2Headers (headers) {
return result return result
} }
/**
* @param {import('node:http2').IncomingHttpHeaders} headers
* @returns {Buffer[]}
*/
function parseH2ResponseHeaders (headers) {
const { [HTTP2_HEADER_STATUS]: _statusCode, ...realHeaders } = headers
return parseH2Headers(realHeaders)
}
async function connectH2 (client, socket) { async function connectH2 (client, socket) {
client[kSocket] = socket client[kSocket] = socket
@@ -19253,22 +19348,32 @@ function writeH2 (client, request) {
headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}` headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}`
headers[HTTP2_HEADER_METHOD] = method headers[HTTP2_HEADER_METHOD] = method
const abort = (err) => { /**
if (request.aborted || request.completed) { * @param {Error} [error]
*/
const abort = (error) => {
if (request.aborted) {
return return
} }
err = err || new RequestAbortedError() if (request.completed) {
if (method === 'CONNECT' && stream != null) {
util.destroy(stream, error || new RequestAbortedError())
}
return
}
util.errorRequest(client, request, err) error = error || new RequestAbortedError()
util.errorRequest(client, request, error)
if (stream != null) { if (stream != null) {
util.destroy(stream, err) util.destroy(stream, error)
} }
// We do not destroy the socket as we can continue using the session // We do not destroy the socket as we can continue using the session
// the stream get's destroyed and the session remains to create new streams // the stream get's destroyed and the session remains to create new streams
util.destroy(body, err) util.destroy(body, error)
client[kQueue][client[kRunningIdx]++] = null client[kQueue][client[kRunningIdx]++] = null
client[kResume]() client[kResume]()
} }
@@ -19287,25 +19392,57 @@ function writeH2 (client, request) {
if (method === 'CONNECT') { if (method === 'CONNECT') {
session.ref() session.ref()
// We are already connected, streams are pending, first request
// will create a new stream. We trigger a request to create the stream and wait until
// `ready` event is triggered
// We disabled endStream to allow the user to write to the stream // We disabled endStream to allow the user to write to the stream
stream = session.request(headers, { endStream: false, signal }) stream = session.request(headers, { endStream: false, signal })
let upgradeResponseFinished = false
if (stream.id && !stream.pending) { /**
request.onUpgrade(null, null, stream) * @param {import('node:http2').IncomingHttpHeaders} headers
++session[kOpenStreams] */
client[kQueue][client[kRunningIdx]++] = null const onResponse = (headers) => {
} else { upgradeResponseFinished = true
stream.once('ready', () => { stream.off(errorMonitor, onUpgradeError)
request.onUpgrade(null, null, stream) request.onUpgradeResponse(Number(headers[HTTP2_HEADER_STATUS]), headers, parseH2ResponseHeaders)
++session[kOpenStreams]
client[kQueue][client[kRunningIdx]++] = null
})
} }
/**
* @param {Error} error
*/
const onUpgradeError = (error) => {
upgradeResponseFinished = true
stream.off('response', onResponse)
request.onUpgradeError(error)
}
const onReady = () => {
try {
request.onUpgrade(null, null, stream)
} catch (error) {
stream.off('response', onResponse)
abort(error)
return
}
if (request.aborted) {
return
}
stream.off('error', abort)
stream.once(errorMonitor, onUpgradeError)
client[kQueue][client[kRunningIdx]++] = null
}
stream.once('response', onResponse)
stream.once('error', abort)
++session[kOpenStreams]
onReady()
stream.once('close', () => { stream.once('close', () => {
if (!upgradeResponseFinished && request.completed) {
stream.off('response', onResponse)
stream.off(errorMonitor, onUpgradeError)
request.onUpgradeError(new InformationalError(`HTTP/2: "stream error" received - code ${stream.rstCode}`))
}
session[kOpenStreams] -= 1 session[kOpenStreams] -= 1
if (session[kOpenStreams] === 0) session.unref() if (session[kOpenStreams] === 0) session.unref()
}) })
@@ -22004,6 +22141,7 @@ class RetryHandler {
this.end = null this.end = null
this.etag = null this.etag = null
this.resume = null this.resume = null
this.headersSent = false
// Handle possible onConnect duplication // Handle possible onConnect duplication
this.handler.onConnect(reason => { this.handler.onConnect(reason => {
@@ -22016,6 +22154,20 @@ class RetryHandler {
}) })
} }
checkpointResponseEnd (headers, resume) {
if (this.end == null && this.opts.method !== 'HEAD') {
const contentLength = headers['content-length']
this.end = contentLength != null ? Number(contentLength) - 1 : null
assert(
this.end == null || Number.isFinite(this.end),
'invalid content-length'
)
}
this.resume = this.end != null ? resume : null
}
onRequestSent () { onRequestSent () {
if (this.handler.onRequestSent) { if (this.handler.onRequestSent) {
this.handler.onRequestSent() this.handler.onRequestSent()
@@ -22104,7 +22256,12 @@ class RetryHandler {
this.retryCount += 1 this.retryCount += 1
if (statusCode >= 300) { if (statusCode >= 300) {
if (this.retryOpts.statusCodes.includes(statusCode) === false) { // Only expose a response if no earlier attempt has reached the caller.
// Otherwise abort this attempt so the error settles the existing body
// instead of replacing it with a new response.
if (!this.headersSent && this.retryOpts.statusCodes.includes(statusCode) === false) {
this.headersSent = true
this.checkpointResponseEnd(headers, resume)
return this.handler.onHeaders( return this.handler.onHeaders(
statusCode, statusCode,
rawHeaders, rawHeaders,
@@ -22173,8 +22330,15 @@ class RetryHandler {
const { start, size, end = size - 1 } = contentRange const { start, size, end = size - 1 } = contentRange
assert(this.start === start, 'content-range mismatch') if (this.start !== start || (this.end != null && this.end !== end)) {
assert(this.end == null || this.end === end, 'content-range mismatch') this.abort(
new RequestRetryError('Content-Range mismatch', statusCode, {
headers,
data: { count: this.retryCount }
})
)
return false
}
this.resume = resume this.resume = resume
return true return true
@@ -22186,6 +22350,7 @@ class RetryHandler {
const range = parseRangeHeader(headers['content-range']) const range = parseRangeHeader(headers['content-range'])
if (range == null) { if (range == null) {
this.headersSent = true
return this.handler.onHeaders( return this.handler.onHeaders(
statusCode, statusCode,
rawHeaders, rawHeaders,
@@ -22224,6 +22389,7 @@ class RetryHandler {
) )
this.resume = resume this.resume = resume
this.headersSent = true
this.etag = headers.etag != null ? headers.etag : null this.etag = headers.etag != null ? headers.etag : null
// Weak etags are not useful for comparison nor cache // Weak etags are not useful for comparison nor cache
@@ -22263,7 +22429,7 @@ class RetryHandler {
} }
onError (err) { onError (err) {
if (this.aborted || isDisturbed(this.opts.body)) { if (this.aborted || isDisturbed(this.opts.body) || (this.headersSent && this.resume == null)) {
return this.handler.onError(err) return this.handler.onError(err)
} }
@@ -26721,6 +26887,49 @@ const COLON = 0x3A
*/ */
const SPACE = 0x20 const SPACE = 0x20
const DATA = Buffer.from('data')
const EVENT = Buffer.from('event')
const ID = Buffer.from('id')
const RETRY = Buffer.from('retry')
function isASCIINumberBytes (buffer, start) {
if (start >= buffer.length) {
return false
}
for (let i = start; i < buffer.length; i++) {
if (buffer[i] < 0x30 || buffer[i] > 0x39) {
return false
}
}
return true
}
function isValidLastEventIdBytes (buffer, start) {
for (let i = start; i < buffer.length; i++) {
if (buffer[i] === 0x00) {
return false
}
}
return true
}
function isFieldName (line, length, field) {
if (length !== field.length) {
return false
}
for (let i = 0; i < length; i++) {
if (line[i] !== field[i]) {
return false
}
}
return true
}
/** /**
* @typedef {object} EventSourceStreamEvent * @typedef {object} EventSourceStreamEvent
* @type {object} * @type {object}
@@ -26761,11 +26970,14 @@ class EventSourceStream extends Transform {
eventEndCheck = false eventEndCheck = false
/** /**
* @type {Buffer} * @type {Buffer[]}
*/ */
buffer = null chunks = []
chunkIndex = 0
pos = 0 pos = 0
lineChunkIndex = 0
linePos = 0
event = { event = {
data: undefined, data: undefined,
@@ -26804,92 +27016,20 @@ class EventSourceStream extends Transform {
return return
} }
// Cache the chunk in the buffer, as the data might not be complete while this.chunks.push(chunk)
// processing it
// TODO: Investigate if there is a more performant way to handle
// incoming chunks
// see: https://github.com/nodejs/undici/issues/2630
if (this.buffer) {
this.buffer = Buffer.concat([this.buffer, chunk])
} else {
this.buffer = chunk
}
// Strip leading byte-order-mark if we opened the stream and started // Strip leading byte-order-mark if we opened the stream and started
// the processing of the incoming data // the processing of the incoming data
if (this.checkBOM) { if (this.checkBOM) {
switch (this.buffer.length) { if (this.handleBOM()) {
case 1: callback()
// Check if the first byte is the same as the first byte of the BOM return
if (this.buffer[0] === BOM[0]) {
// If it is, we need to wait for more data
callback()
return
}
// Set the checkBOM flag to false as we don't need to check for the
// BOM anymore
this.checkBOM = false
// The buffer only contains one byte so we need to wait for more data
callback()
return
case 2:
// Check if the first two bytes are the same as the first two bytes
// of the BOM
if (
this.buffer[0] === BOM[0] &&
this.buffer[1] === BOM[1]
) {
// If it is, we need to wait for more data, because the third byte
// is needed to determine if it is the BOM or not
callback()
return
}
// Set the checkBOM flag to false as we don't need to check for the
// BOM anymore
this.checkBOM = false
break
case 3:
// Check if the first three bytes are the same as the first three
// bytes of the BOM
if (
this.buffer[0] === BOM[0] &&
this.buffer[1] === BOM[1] &&
this.buffer[2] === BOM[2]
) {
// If it is, we can drop the buffered data, as it is only the BOM
this.buffer = Buffer.alloc(0)
// Set the checkBOM flag to false as we don't need to check for the
// BOM anymore
this.checkBOM = false
// Await more data
callback()
return
}
// If it is not the BOM, we can start processing the data
this.checkBOM = false
break
default:
// The buffer is longer than 3 bytes, so we can drop the BOM if it is
// present
if (
this.buffer[0] === BOM[0] &&
this.buffer[1] === BOM[1] &&
this.buffer[2] === BOM[2]
) {
// Remove the BOM from the buffer
this.buffer = this.buffer.subarray(3)
}
// Set the checkBOM flag to false as we don't need to check for the
this.checkBOM = false
break
} }
} }
while (this.pos < this.buffer.length) { while (this.hasCurrentByte()) {
const byte = this.currentByte()
// If the previous line ended with an end-of-line, we need to check // If the previous line ended with an end-of-line, we need to check
// if the next character is also an end-of-line. // if the next character is also an end-of-line.
if (this.eventEndCheck) { if (this.eventEndCheck) {
@@ -26902,10 +27042,9 @@ class EventSourceStream extends Transform {
if (this.crlfCheck) { if (this.crlfCheck) {
// If the current character is a line feed, we can remove it // If the current character is a line feed, we can remove it
// from the buffer and reset the crlfCheck flag // from the buffer and reset the crlfCheck flag
if (this.buffer[this.pos] === LF) { if (byte === LF) {
this.buffer = this.buffer.subarray(this.pos + 1)
this.pos = 0
this.crlfCheck = false this.crlfCheck = false
this.consumeCurrentByte()
// It is possible that the line feed is not the end of the // It is possible that the line feed is not the end of the
// event. We need to check if the next character is an // event. We need to check if the next character is an
@@ -26921,19 +27060,17 @@ class EventSourceStream extends Transform {
this.crlfCheck = false this.crlfCheck = false
} }
if (this.buffer[this.pos] === LF || this.buffer[this.pos] === CR) { if (byte === LF || byte === CR) {
// If the current character is a carriage return, we need to // If the current character is a carriage return, we need to
// set the crlfCheck flag to true, as we need to check if the // set the crlfCheck flag to true, as we need to check if the
// next character is a line feed so we can remove it from the // next character is a line feed so we can remove it from the
// buffer // buffer
if (this.buffer[this.pos] === CR) { if (byte === CR) {
this.crlfCheck = true this.crlfCheck = true
} }
this.buffer = this.buffer.subarray(this.pos + 1) this.consumeCurrentByte()
this.pos = 0 if (this.hasPendingEvent()) {
if (
this.event.data !== undefined || this.event.event || this.event.id || this.event.retry) {
this.processEvent(this.event) this.processEvent(this.event)
} }
this.clearEvent() this.clearEvent()
@@ -26947,22 +27084,18 @@ class EventSourceStream extends Transform {
// If the current character is an end-of-line, we can process the // If the current character is an end-of-line, we can process the
// line // line
if (this.buffer[this.pos] === LF || this.buffer[this.pos] === CR) { if (byte === LF || byte === CR) {
// If the current character is a carriage return, we need to // If the current character is a carriage return, we need to
// set the crlfCheck flag to true, as we need to check if the // set the crlfCheck flag to true, as we need to check if the
// next character is a line feed // next character is a line feed
if (this.buffer[this.pos] === CR) { if (byte === CR) {
this.crlfCheck = true this.crlfCheck = true
} }
// In any case, we can process the line as we reached an // In any case, we can process the line as we reached an
// end-of-line character // end-of-line character
this.parseLine(this.buffer.subarray(0, this.pos), this.event) this.parseLine(this.readLine(), this.event)
this.consumeCurrentByte()
// Remove the processed line from the buffer
this.buffer = this.buffer.subarray(this.pos + 1)
// Reset the position as we removed the processed line from the buffer
this.pos = 0
// A line was processed and this could be the end of the event. We need // A line was processed and this could be the end of the event. We need
// to check if the next line is empty to determine if the event is // to check if the next line is empty to determine if the event is
// finished. // finished.
@@ -26970,7 +27103,7 @@ class EventSourceStream extends Transform {
continue continue
} }
this.pos++ this.advanceCursor()
} }
callback() callback()
@@ -26995,64 +27128,53 @@ class EventSourceStream extends Transform {
return return
} }
let field = '' let fieldLength = line.length
let value = '' let valueStart = line.length
// If the line contains a U+003A COLON character (:) // If the line contains a U+003A COLON character (:)
if (colonPosition !== -1) { if (colonPosition !== -1) {
// Collect the characters on the line before the first U+003A COLON fieldLength = colonPosition
// character (:), and let field be that string.
// TODO: Investigate if there is a more performant way to extract the
// field
// see: https://github.com/nodejs/undici/issues/2630
field = line.subarray(0, colonPosition).toString('utf8')
// Collect the characters on the line after the first U+003A COLON // Collect the characters on the line after the first U+003A COLON
// character (:), and let value be that string. // character (:), and let value be that string.
// If value starts with a U+0020 SPACE character, remove it from value. // If value starts with a U+0020 SPACE character, remove it from value.
let valueStart = colonPosition + 1 valueStart = colonPosition + 1
if (line[valueStart] === SPACE) { if (line[valueStart] === SPACE) {
++valueStart ++valueStart
} }
// TODO: Investigate if there is a more performant way to extract the
// value
// see: https://github.com/nodejs/undici/issues/2630
value = line.subarray(valueStart).toString('utf8')
// Otherwise, the string is not empty but does not contain a U+003A COLON
// character (:)
} else {
// Process the field using the steps described below, using the whole
// line as the field name, and the empty string as the field value.
field = line.toString('utf8')
value = ''
} }
// Modify the event with the field name and value. The value is also if (isFieldName(line, fieldLength, DATA)) {
// decoded as UTF-8 const value = line.toString('utf8', valueStart)
switch (field) {
case 'data': if (event.data === undefined) {
if (event[field] === undefined) { event.data = value
event[field] = value } else {
} else { event.data += `\n${value}`
event[field] += `\n${value}` }
} return
break }
case 'retry':
if (isASCIINumber(value)) { if (isFieldName(line, fieldLength, RETRY)) {
event[field] = value if (isASCIINumberBytes(line, valueStart)) {
} event.retry = line.toString('utf8', valueStart)
break }
case 'id': return
if (isValidLastEventId(value)) { }
event[field] = value
} if (isFieldName(line, fieldLength, ID)) {
break if (isValidLastEventIdBytes(line, valueStart)) {
case 'event': event.id = line.toString('utf8', valueStart)
if (value.length > 0) { }
event[field] = value return
} }
break
if (isFieldName(line, fieldLength, EVENT)) {
const value = line.toString('utf8', valueStart)
if (value.length > 0) {
event.event = value
}
} }
} }
@@ -27082,13 +27204,152 @@ class EventSourceStream extends Transform {
} }
clearEvent () { clearEvent () {
this.event = { this.event.data = undefined
data: undefined, this.event.event = undefined
event: undefined, this.event.id = undefined
id: undefined, this.event.retry = undefined
retry: undefined }
hasPendingEvent () {
return this.event.data !== undefined ||
this.event.event !== undefined ||
this.event.id !== undefined ||
this.event.retry !== undefined
}
hasCurrentByte () {
return this.chunkIndex < this.chunks.length &&
this.pos < this.chunks[this.chunkIndex].length
}
currentByte () {
return this.chunks[this.chunkIndex][this.pos]
}
consumeCurrentByte () {
this.advanceCursor()
this.syncLineStartToCursor()
}
advanceCursor () {
this.pos++
while (this.chunkIndex < this.chunks.length && this.pos >= this.chunks[this.chunkIndex].length) {
this.chunkIndex++
this.pos = 0
} }
} }
syncLineStartToCursor () {
this.lineChunkIndex = this.chunkIndex
this.linePos = this.pos
this.dropConsumedChunks()
}
dropConsumedChunks () {
while (this.lineChunkIndex > 0) {
this.chunks.shift()
this.lineChunkIndex--
this.chunkIndex--
}
if (this.chunkIndex === this.chunks.length) {
this.chunks.length = 0
this.chunkIndex = 0
this.pos = 0
this.lineChunkIndex = 0
this.linePos = 0
}
}
readLine () {
if (this.lineChunkIndex === this.chunkIndex) {
return this.chunks[this.chunkIndex].subarray(this.linePos, this.pos)
}
const chunks = []
let length = 0
for (let i = this.lineChunkIndex; i <= this.chunkIndex; i++) {
const chunk = this.chunks[i]
const start = i === this.lineChunkIndex ? this.linePos : 0
const end = i === this.chunkIndex ? this.pos : chunk.length
const slice = chunk.subarray(start, end)
length += slice.length
chunks.push(slice)
}
return Buffer.concat(chunks, length)
}
peekBufferedByte (offset) {
let chunkIndex = this.lineChunkIndex
let pos = this.linePos
while (chunkIndex < this.chunks.length) {
const chunk = this.chunks[chunkIndex]
const remaining = chunk.length - pos
if (offset < remaining) {
return chunk[pos + offset]
}
offset -= remaining
chunkIndex++
pos = 0
}
}
discardLeadingBytes (count) {
while (count > 0 && this.lineChunkIndex < this.chunks.length) {
const chunk = this.chunks[this.lineChunkIndex]
const remaining = chunk.length - this.linePos
if (count < remaining) {
this.linePos += count
count = 0
} else {
count -= remaining
this.lineChunkIndex++
this.linePos = 0
}
}
this.chunkIndex = this.lineChunkIndex
this.pos = this.linePos
this.dropConsumedChunks()
}
handleBOM () {
const first = this.peekBufferedByte(0)
const second = this.peekBufferedByte(1)
const third = this.peekBufferedByte(2)
if (second === undefined) {
if (first === BOM[0]) {
return true
}
this.checkBOM = false
return true
}
if (third === undefined) {
if (first === BOM[0] && second === BOM[1]) {
return true
}
this.checkBOM = false
return false
}
if (first === BOM[0] && second === BOM[1] && third === BOM[2]) {
this.discardLeadingBytes(3)
}
this.checkBOM = false
return !this.hasCurrentByte()
}
} }
module.exports = { module.exports = {
@@ -38356,7 +38617,7 @@ function establishWebSocketConnection (url, protocols, client, ws, onEstablish,
// is specified, the server needs to include the same field and one of // is specified, the server needs to include the same field and one of
// the selected subprotocol values in its response for the connection to // the selected subprotocol values in its response for the connection to
// be established. // be established.
if (!requestProtocols.includes(secProtocol)) { if (requestProtocols === null || !requestProtocols.includes(secProtocol)) {
failWebsocketConnection(ws, 'Protocol was not set in the opening handshake.') failWebsocketConnection(ws, 'Protocol was not set in the opening handshake.')
return return
} }
@@ -39117,7 +39378,12 @@ class PerMessageDeflate {
if (this.#maxPayloadSize > 0 && this.#inflate[kLength] > this.#maxPayloadSize) { if (this.#maxPayloadSize > 0 && this.#inflate[kLength] > this.#maxPayloadSize) {
callback(new MessageSizeExceededError()) callback(new MessageSizeExceededError())
// The inflater may still hold buffered input that can emit a late
// zlib error. Remove the data listener, then deterministically stop
// the stream so a subsequent 'error' cannot fire without a listener
// (which would terminate the process as an unhandled error event).
this.#inflate.removeAllListeners() this.#inflate.removeAllListeners()
this.#inflate.destroy()
this.#inflate = null this.#inflate = null
return return
} }

640
dist/setup/index.js vendored
View File

@@ -15982,11 +15982,77 @@ class Request {
} }
} }
onUpgrade (statusCode, headers, socket) { /**
* @param {number|null} statusCode
* @param {Buffer[]|null} headers
* @param {import('node:stream').Duplex} socket
* @param {string} [statusText]
*/
onUpgrade (statusCode, headers, socket, statusText = '') {
this.onFinally()
assert(!this.aborted) assert(!this.aborted)
assert(!this.completed) assert(!this.completed)
return this[kHandler].onUpgrade(statusCode, headers, socket) if (statusCode !== null) {
this.#publishUpgradeHeaders(statusCode, headers, statusText)
}
const result = this[kHandler].onUpgrade(statusCode, headers, socket)
if (!this.aborted) {
this.completed = true
if (statusCode !== null) {
this.#publishUpgradeTrailers()
}
}
return result
}
/**
* @param {number} statusCode
* @param {import('node:http2').IncomingHttpHeaders} headers
* @param {(headers: import('node:http2').IncomingHttpHeaders) => Buffer[]} parseHeaders
* @param {string} [statusText]
*/
onUpgradeResponse (statusCode, headers, parseHeaders, statusText = '') {
assert(!this.aborted)
assert(this.completed)
if (channels.headers.hasSubscribers) {
this.#publishUpgradeHeaders(statusCode, parseHeaders(headers), statusText)
}
this.#publishUpgradeTrailers()
}
/**
* @param {Error} error
*/
onUpgradeError (error) {
assert(!this.aborted)
assert(this.completed)
if (channels.error.hasSubscribers) {
channels.error.publish({ request: this, error })
}
}
/**
* @param {number} statusCode
* @param {Buffer[]} headers
* @param {string} statusText
*/
#publishUpgradeHeaders (statusCode, headers, statusText) {
if (channels.headers.hasSubscribers) {
channels.headers.publish({ request: this, response: { statusCode, headers, statusText } })
}
}
#publishUpgradeTrailers () {
if (channels.trailers.hasSubscribers) {
channels.trailers.publish({ request: this, trailers: [] })
}
} }
onComplete (trailers) { onComplete (trailers) {
@@ -17885,7 +17951,7 @@ class Parser {
} }
onUpgrade (head) { onUpgrade (head) {
const { upgrade, client, socket, headers, statusCode } = this const { upgrade, client, socket, headers, statusCode, statusText } = this
assert(upgrade) assert(upgrade)
assert(client[kSocket] === socket) assert(client[kSocket] === socket)
@@ -17920,9 +17986,10 @@ class Parser {
client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade')) client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade'))
try { try {
request.onUpgrade(statusCode, headers, socket) request.onUpgrade(statusCode, headers, socket, statusText)
} catch (err) { } catch (error) {
util.destroy(socket, err) util.errorRequest(client, request, error)
util.destroy(socket, error)
} }
client[kResume]() client[kResume]()
@@ -18329,7 +18396,7 @@ async function connectH1 (client, socket) {
function clearIdleSocketValidation (socket) { function clearIdleSocketValidation (socket) {
if (socket[kIdleSocketValidationTimeout]) { if (socket[kIdleSocketValidationTimeout]) {
clearTimeout(socket[kIdleSocketValidationTimeout]) clearImmediate(socket[kIdleSocketValidationTimeout])
socket[kIdleSocketValidationTimeout] = null socket[kIdleSocketValidationTimeout] = null
} }
@@ -18338,15 +18405,23 @@ function clearIdleSocketValidation (socket) {
function scheduleIdleSocketValidation (client, socket) { function scheduleIdleSocketValidation (client, socket) {
socket[kIdleSocketValidation] = 1 socket[kIdleSocketValidation] = 1
socket[kIdleSocketValidationTimeout] = setTimeout(() => { // Yield to the check phase (after poll) so unsolicited bytes / FIN / RST
// already pending on this idle keep-alive socket are processed before the
// next request is written (GHSA-35p6-xmwp-9g52).
//
// setTimeout(0) pays Node's ~1ms timer floor on every sequential reuse
// (#5493). setImmediate avoids that, but an *unref'd* Immediate lets poll
// block for ~500ms when the event loop is otherwise idle (#5600 / #5606).
// A ref'd Immediate both keeps the pending request alive and makes poll
// return immediately — the hybrid those issues asked for.
socket[kIdleSocketValidationTimeout] = setImmediate(() => {
socket[kIdleSocketValidationTimeout] = null socket[kIdleSocketValidationTimeout] = null
socket[kIdleSocketValidation] = 2 socket[kIdleSocketValidation] = 2
if (client[kSocket] === socket && !socket.destroyed) { if (client[kSocket] === socket && !socket.destroyed) {
client[kResume]() client[kResume]()
} }
}, 0) })
socket[kIdleSocketValidationTimeout].unref?.()
} }
/** /**
@@ -18495,12 +18570,22 @@ function writeH1 (client, request) {
const socket = client[kSocket] const socket = client[kSocket]
clearIdleSocketValidation(socket) clearIdleSocketValidation(socket)
const abort = (err) => { /**
if (request.aborted || request.completed) { * @param {Error} [error]
*/
const abort = (error) => {
if (request.aborted) {
return return
} }
util.errorRequest(client, request, err || new RequestAbortedError()) if (request.completed) {
if (request.upgrade || request.method === 'CONNECT') {
util.destroy(socket, new InformationalError('aborted'))
}
return
}
util.errorRequest(client, request, error || new RequestAbortedError())
util.destroy(body) util.destroy(body)
util.destroy(socket, new InformationalError('aborted')) util.destroy(socket, new InformationalError('aborted'))
@@ -18957,6 +19042,7 @@ module.exports = connectH1
const assert = __nccwpck_require__(4589) const assert = __nccwpck_require__(4589)
const { errorMonitor } = __nccwpck_require__(8474)
const { pipeline } = __nccwpck_require__(7075) const { pipeline } = __nccwpck_require__(7075)
const util = __nccwpck_require__(3440) const util = __nccwpck_require__(3440)
const { const {
@@ -19033,6 +19119,15 @@ function parseH2Headers (headers) {
return result return result
} }
/**
* @param {import('node:http2').IncomingHttpHeaders} headers
* @returns {Buffer[]}
*/
function parseH2ResponseHeaders (headers) {
const { [HTTP2_HEADER_STATUS]: _statusCode, ...realHeaders } = headers
return parseH2Headers(realHeaders)
}
async function connectH2 (client, socket) { async function connectH2 (client, socket) {
client[kSocket] = socket client[kSocket] = socket
@@ -19253,22 +19348,32 @@ function writeH2 (client, request) {
headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}` headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}`
headers[HTTP2_HEADER_METHOD] = method headers[HTTP2_HEADER_METHOD] = method
const abort = (err) => { /**
if (request.aborted || request.completed) { * @param {Error} [error]
*/
const abort = (error) => {
if (request.aborted) {
return return
} }
err = err || new RequestAbortedError() if (request.completed) {
if (method === 'CONNECT' && stream != null) {
util.destroy(stream, error || new RequestAbortedError())
}
return
}
util.errorRequest(client, request, err) error = error || new RequestAbortedError()
util.errorRequest(client, request, error)
if (stream != null) { if (stream != null) {
util.destroy(stream, err) util.destroy(stream, error)
} }
// We do not destroy the socket as we can continue using the session // We do not destroy the socket as we can continue using the session
// the stream get's destroyed and the session remains to create new streams // the stream get's destroyed and the session remains to create new streams
util.destroy(body, err) util.destroy(body, error)
client[kQueue][client[kRunningIdx]++] = null client[kQueue][client[kRunningIdx]++] = null
client[kResume]() client[kResume]()
} }
@@ -19287,25 +19392,57 @@ function writeH2 (client, request) {
if (method === 'CONNECT') { if (method === 'CONNECT') {
session.ref() session.ref()
// We are already connected, streams are pending, first request
// will create a new stream. We trigger a request to create the stream and wait until
// `ready` event is triggered
// We disabled endStream to allow the user to write to the stream // We disabled endStream to allow the user to write to the stream
stream = session.request(headers, { endStream: false, signal }) stream = session.request(headers, { endStream: false, signal })
let upgradeResponseFinished = false
if (stream.id && !stream.pending) { /**
request.onUpgrade(null, null, stream) * @param {import('node:http2').IncomingHttpHeaders} headers
++session[kOpenStreams] */
client[kQueue][client[kRunningIdx]++] = null const onResponse = (headers) => {
} else { upgradeResponseFinished = true
stream.once('ready', () => { stream.off(errorMonitor, onUpgradeError)
request.onUpgrade(null, null, stream) request.onUpgradeResponse(Number(headers[HTTP2_HEADER_STATUS]), headers, parseH2ResponseHeaders)
++session[kOpenStreams]
client[kQueue][client[kRunningIdx]++] = null
})
} }
/**
* @param {Error} error
*/
const onUpgradeError = (error) => {
upgradeResponseFinished = true
stream.off('response', onResponse)
request.onUpgradeError(error)
}
const onReady = () => {
try {
request.onUpgrade(null, null, stream)
} catch (error) {
stream.off('response', onResponse)
abort(error)
return
}
if (request.aborted) {
return
}
stream.off('error', abort)
stream.once(errorMonitor, onUpgradeError)
client[kQueue][client[kRunningIdx]++] = null
}
stream.once('response', onResponse)
stream.once('error', abort)
++session[kOpenStreams]
onReady()
stream.once('close', () => { stream.once('close', () => {
if (!upgradeResponseFinished && request.completed) {
stream.off('response', onResponse)
stream.off(errorMonitor, onUpgradeError)
request.onUpgradeError(new InformationalError(`HTTP/2: "stream error" received - code ${stream.rstCode}`))
}
session[kOpenStreams] -= 1 session[kOpenStreams] -= 1
if (session[kOpenStreams] === 0) session.unref() if (session[kOpenStreams] === 0) session.unref()
}) })
@@ -22004,6 +22141,7 @@ class RetryHandler {
this.end = null this.end = null
this.etag = null this.etag = null
this.resume = null this.resume = null
this.headersSent = false
// Handle possible onConnect duplication // Handle possible onConnect duplication
this.handler.onConnect(reason => { this.handler.onConnect(reason => {
@@ -22016,6 +22154,20 @@ class RetryHandler {
}) })
} }
checkpointResponseEnd (headers, resume) {
if (this.end == null && this.opts.method !== 'HEAD') {
const contentLength = headers['content-length']
this.end = contentLength != null ? Number(contentLength) - 1 : null
assert(
this.end == null || Number.isFinite(this.end),
'invalid content-length'
)
}
this.resume = this.end != null ? resume : null
}
onRequestSent () { onRequestSent () {
if (this.handler.onRequestSent) { if (this.handler.onRequestSent) {
this.handler.onRequestSent() this.handler.onRequestSent()
@@ -22104,7 +22256,12 @@ class RetryHandler {
this.retryCount += 1 this.retryCount += 1
if (statusCode >= 300) { if (statusCode >= 300) {
if (this.retryOpts.statusCodes.includes(statusCode) === false) { // Only expose a response if no earlier attempt has reached the caller.
// Otherwise abort this attempt so the error settles the existing body
// instead of replacing it with a new response.
if (!this.headersSent && this.retryOpts.statusCodes.includes(statusCode) === false) {
this.headersSent = true
this.checkpointResponseEnd(headers, resume)
return this.handler.onHeaders( return this.handler.onHeaders(
statusCode, statusCode,
rawHeaders, rawHeaders,
@@ -22173,8 +22330,15 @@ class RetryHandler {
const { start, size, end = size - 1 } = contentRange const { start, size, end = size - 1 } = contentRange
assert(this.start === start, 'content-range mismatch') if (this.start !== start || (this.end != null && this.end !== end)) {
assert(this.end == null || this.end === end, 'content-range mismatch') this.abort(
new RequestRetryError('Content-Range mismatch', statusCode, {
headers,
data: { count: this.retryCount }
})
)
return false
}
this.resume = resume this.resume = resume
return true return true
@@ -22186,6 +22350,7 @@ class RetryHandler {
const range = parseRangeHeader(headers['content-range']) const range = parseRangeHeader(headers['content-range'])
if (range == null) { if (range == null) {
this.headersSent = true
return this.handler.onHeaders( return this.handler.onHeaders(
statusCode, statusCode,
rawHeaders, rawHeaders,
@@ -22224,6 +22389,7 @@ class RetryHandler {
) )
this.resume = resume this.resume = resume
this.headersSent = true
this.etag = headers.etag != null ? headers.etag : null this.etag = headers.etag != null ? headers.etag : null
// Weak etags are not useful for comparison nor cache // Weak etags are not useful for comparison nor cache
@@ -22263,7 +22429,7 @@ class RetryHandler {
} }
onError (err) { onError (err) {
if (this.aborted || isDisturbed(this.opts.body)) { if (this.aborted || isDisturbed(this.opts.body) || (this.headersSent && this.resume == null)) {
return this.handler.onError(err) return this.handler.onError(err)
} }
@@ -26721,6 +26887,49 @@ const COLON = 0x3A
*/ */
const SPACE = 0x20 const SPACE = 0x20
const DATA = Buffer.from('data')
const EVENT = Buffer.from('event')
const ID = Buffer.from('id')
const RETRY = Buffer.from('retry')
function isASCIINumberBytes (buffer, start) {
if (start >= buffer.length) {
return false
}
for (let i = start; i < buffer.length; i++) {
if (buffer[i] < 0x30 || buffer[i] > 0x39) {
return false
}
}
return true
}
function isValidLastEventIdBytes (buffer, start) {
for (let i = start; i < buffer.length; i++) {
if (buffer[i] === 0x00) {
return false
}
}
return true
}
function isFieldName (line, length, field) {
if (length !== field.length) {
return false
}
for (let i = 0; i < length; i++) {
if (line[i] !== field[i]) {
return false
}
}
return true
}
/** /**
* @typedef {object} EventSourceStreamEvent * @typedef {object} EventSourceStreamEvent
* @type {object} * @type {object}
@@ -26761,11 +26970,14 @@ class EventSourceStream extends Transform {
eventEndCheck = false eventEndCheck = false
/** /**
* @type {Buffer} * @type {Buffer[]}
*/ */
buffer = null chunks = []
chunkIndex = 0
pos = 0 pos = 0
lineChunkIndex = 0
linePos = 0
event = { event = {
data: undefined, data: undefined,
@@ -26804,92 +27016,20 @@ class EventSourceStream extends Transform {
return return
} }
// Cache the chunk in the buffer, as the data might not be complete while this.chunks.push(chunk)
// processing it
// TODO: Investigate if there is a more performant way to handle
// incoming chunks
// see: https://github.com/nodejs/undici/issues/2630
if (this.buffer) {
this.buffer = Buffer.concat([this.buffer, chunk])
} else {
this.buffer = chunk
}
// Strip leading byte-order-mark if we opened the stream and started // Strip leading byte-order-mark if we opened the stream and started
// the processing of the incoming data // the processing of the incoming data
if (this.checkBOM) { if (this.checkBOM) {
switch (this.buffer.length) { if (this.handleBOM()) {
case 1: callback()
// Check if the first byte is the same as the first byte of the BOM return
if (this.buffer[0] === BOM[0]) {
// If it is, we need to wait for more data
callback()
return
}
// Set the checkBOM flag to false as we don't need to check for the
// BOM anymore
this.checkBOM = false
// The buffer only contains one byte so we need to wait for more data
callback()
return
case 2:
// Check if the first two bytes are the same as the first two bytes
// of the BOM
if (
this.buffer[0] === BOM[0] &&
this.buffer[1] === BOM[1]
) {
// If it is, we need to wait for more data, because the third byte
// is needed to determine if it is the BOM or not
callback()
return
}
// Set the checkBOM flag to false as we don't need to check for the
// BOM anymore
this.checkBOM = false
break
case 3:
// Check if the first three bytes are the same as the first three
// bytes of the BOM
if (
this.buffer[0] === BOM[0] &&
this.buffer[1] === BOM[1] &&
this.buffer[2] === BOM[2]
) {
// If it is, we can drop the buffered data, as it is only the BOM
this.buffer = Buffer.alloc(0)
// Set the checkBOM flag to false as we don't need to check for the
// BOM anymore
this.checkBOM = false
// Await more data
callback()
return
}
// If it is not the BOM, we can start processing the data
this.checkBOM = false
break
default:
// The buffer is longer than 3 bytes, so we can drop the BOM if it is
// present
if (
this.buffer[0] === BOM[0] &&
this.buffer[1] === BOM[1] &&
this.buffer[2] === BOM[2]
) {
// Remove the BOM from the buffer
this.buffer = this.buffer.subarray(3)
}
// Set the checkBOM flag to false as we don't need to check for the
this.checkBOM = false
break
} }
} }
while (this.pos < this.buffer.length) { while (this.hasCurrentByte()) {
const byte = this.currentByte()
// If the previous line ended with an end-of-line, we need to check // If the previous line ended with an end-of-line, we need to check
// if the next character is also an end-of-line. // if the next character is also an end-of-line.
if (this.eventEndCheck) { if (this.eventEndCheck) {
@@ -26902,10 +27042,9 @@ class EventSourceStream extends Transform {
if (this.crlfCheck) { if (this.crlfCheck) {
// If the current character is a line feed, we can remove it // If the current character is a line feed, we can remove it
// from the buffer and reset the crlfCheck flag // from the buffer and reset the crlfCheck flag
if (this.buffer[this.pos] === LF) { if (byte === LF) {
this.buffer = this.buffer.subarray(this.pos + 1)
this.pos = 0
this.crlfCheck = false this.crlfCheck = false
this.consumeCurrentByte()
// It is possible that the line feed is not the end of the // It is possible that the line feed is not the end of the
// event. We need to check if the next character is an // event. We need to check if the next character is an
@@ -26921,19 +27060,17 @@ class EventSourceStream extends Transform {
this.crlfCheck = false this.crlfCheck = false
} }
if (this.buffer[this.pos] === LF || this.buffer[this.pos] === CR) { if (byte === LF || byte === CR) {
// If the current character is a carriage return, we need to // If the current character is a carriage return, we need to
// set the crlfCheck flag to true, as we need to check if the // set the crlfCheck flag to true, as we need to check if the
// next character is a line feed so we can remove it from the // next character is a line feed so we can remove it from the
// buffer // buffer
if (this.buffer[this.pos] === CR) { if (byte === CR) {
this.crlfCheck = true this.crlfCheck = true
} }
this.buffer = this.buffer.subarray(this.pos + 1) this.consumeCurrentByte()
this.pos = 0 if (this.hasPendingEvent()) {
if (
this.event.data !== undefined || this.event.event || this.event.id || this.event.retry) {
this.processEvent(this.event) this.processEvent(this.event)
} }
this.clearEvent() this.clearEvent()
@@ -26947,22 +27084,18 @@ class EventSourceStream extends Transform {
// If the current character is an end-of-line, we can process the // If the current character is an end-of-line, we can process the
// line // line
if (this.buffer[this.pos] === LF || this.buffer[this.pos] === CR) { if (byte === LF || byte === CR) {
// If the current character is a carriage return, we need to // If the current character is a carriage return, we need to
// set the crlfCheck flag to true, as we need to check if the // set the crlfCheck flag to true, as we need to check if the
// next character is a line feed // next character is a line feed
if (this.buffer[this.pos] === CR) { if (byte === CR) {
this.crlfCheck = true this.crlfCheck = true
} }
// In any case, we can process the line as we reached an // In any case, we can process the line as we reached an
// end-of-line character // end-of-line character
this.parseLine(this.buffer.subarray(0, this.pos), this.event) this.parseLine(this.readLine(), this.event)
this.consumeCurrentByte()
// Remove the processed line from the buffer
this.buffer = this.buffer.subarray(this.pos + 1)
// Reset the position as we removed the processed line from the buffer
this.pos = 0
// A line was processed and this could be the end of the event. We need // A line was processed and this could be the end of the event. We need
// to check if the next line is empty to determine if the event is // to check if the next line is empty to determine if the event is
// finished. // finished.
@@ -26970,7 +27103,7 @@ class EventSourceStream extends Transform {
continue continue
} }
this.pos++ this.advanceCursor()
} }
callback() callback()
@@ -26995,64 +27128,53 @@ class EventSourceStream extends Transform {
return return
} }
let field = '' let fieldLength = line.length
let value = '' let valueStart = line.length
// If the line contains a U+003A COLON character (:) // If the line contains a U+003A COLON character (:)
if (colonPosition !== -1) { if (colonPosition !== -1) {
// Collect the characters on the line before the first U+003A COLON fieldLength = colonPosition
// character (:), and let field be that string.
// TODO: Investigate if there is a more performant way to extract the
// field
// see: https://github.com/nodejs/undici/issues/2630
field = line.subarray(0, colonPosition).toString('utf8')
// Collect the characters on the line after the first U+003A COLON // Collect the characters on the line after the first U+003A COLON
// character (:), and let value be that string. // character (:), and let value be that string.
// If value starts with a U+0020 SPACE character, remove it from value. // If value starts with a U+0020 SPACE character, remove it from value.
let valueStart = colonPosition + 1 valueStart = colonPosition + 1
if (line[valueStart] === SPACE) { if (line[valueStart] === SPACE) {
++valueStart ++valueStart
} }
// TODO: Investigate if there is a more performant way to extract the
// value
// see: https://github.com/nodejs/undici/issues/2630
value = line.subarray(valueStart).toString('utf8')
// Otherwise, the string is not empty but does not contain a U+003A COLON
// character (:)
} else {
// Process the field using the steps described below, using the whole
// line as the field name, and the empty string as the field value.
field = line.toString('utf8')
value = ''
} }
// Modify the event with the field name and value. The value is also if (isFieldName(line, fieldLength, DATA)) {
// decoded as UTF-8 const value = line.toString('utf8', valueStart)
switch (field) {
case 'data': if (event.data === undefined) {
if (event[field] === undefined) { event.data = value
event[field] = value } else {
} else { event.data += `\n${value}`
event[field] += `\n${value}` }
} return
break }
case 'retry':
if (isASCIINumber(value)) { if (isFieldName(line, fieldLength, RETRY)) {
event[field] = value if (isASCIINumberBytes(line, valueStart)) {
} event.retry = line.toString('utf8', valueStart)
break }
case 'id': return
if (isValidLastEventId(value)) { }
event[field] = value
} if (isFieldName(line, fieldLength, ID)) {
break if (isValidLastEventIdBytes(line, valueStart)) {
case 'event': event.id = line.toString('utf8', valueStart)
if (value.length > 0) { }
event[field] = value return
} }
break
if (isFieldName(line, fieldLength, EVENT)) {
const value = line.toString('utf8', valueStart)
if (value.length > 0) {
event.event = value
}
} }
} }
@@ -27082,13 +27204,152 @@ class EventSourceStream extends Transform {
} }
clearEvent () { clearEvent () {
this.event = { this.event.data = undefined
data: undefined, this.event.event = undefined
event: undefined, this.event.id = undefined
id: undefined, this.event.retry = undefined
retry: undefined }
hasPendingEvent () {
return this.event.data !== undefined ||
this.event.event !== undefined ||
this.event.id !== undefined ||
this.event.retry !== undefined
}
hasCurrentByte () {
return this.chunkIndex < this.chunks.length &&
this.pos < this.chunks[this.chunkIndex].length
}
currentByte () {
return this.chunks[this.chunkIndex][this.pos]
}
consumeCurrentByte () {
this.advanceCursor()
this.syncLineStartToCursor()
}
advanceCursor () {
this.pos++
while (this.chunkIndex < this.chunks.length && this.pos >= this.chunks[this.chunkIndex].length) {
this.chunkIndex++
this.pos = 0
} }
} }
syncLineStartToCursor () {
this.lineChunkIndex = this.chunkIndex
this.linePos = this.pos
this.dropConsumedChunks()
}
dropConsumedChunks () {
while (this.lineChunkIndex > 0) {
this.chunks.shift()
this.lineChunkIndex--
this.chunkIndex--
}
if (this.chunkIndex === this.chunks.length) {
this.chunks.length = 0
this.chunkIndex = 0
this.pos = 0
this.lineChunkIndex = 0
this.linePos = 0
}
}
readLine () {
if (this.lineChunkIndex === this.chunkIndex) {
return this.chunks[this.chunkIndex].subarray(this.linePos, this.pos)
}
const chunks = []
let length = 0
for (let i = this.lineChunkIndex; i <= this.chunkIndex; i++) {
const chunk = this.chunks[i]
const start = i === this.lineChunkIndex ? this.linePos : 0
const end = i === this.chunkIndex ? this.pos : chunk.length
const slice = chunk.subarray(start, end)
length += slice.length
chunks.push(slice)
}
return Buffer.concat(chunks, length)
}
peekBufferedByte (offset) {
let chunkIndex = this.lineChunkIndex
let pos = this.linePos
while (chunkIndex < this.chunks.length) {
const chunk = this.chunks[chunkIndex]
const remaining = chunk.length - pos
if (offset < remaining) {
return chunk[pos + offset]
}
offset -= remaining
chunkIndex++
pos = 0
}
}
discardLeadingBytes (count) {
while (count > 0 && this.lineChunkIndex < this.chunks.length) {
const chunk = this.chunks[this.lineChunkIndex]
const remaining = chunk.length - this.linePos
if (count < remaining) {
this.linePos += count
count = 0
} else {
count -= remaining
this.lineChunkIndex++
this.linePos = 0
}
}
this.chunkIndex = this.lineChunkIndex
this.pos = this.linePos
this.dropConsumedChunks()
}
handleBOM () {
const first = this.peekBufferedByte(0)
const second = this.peekBufferedByte(1)
const third = this.peekBufferedByte(2)
if (second === undefined) {
if (first === BOM[0]) {
return true
}
this.checkBOM = false
return true
}
if (third === undefined) {
if (first === BOM[0] && second === BOM[1]) {
return true
}
this.checkBOM = false
return false
}
if (first === BOM[0] && second === BOM[1] && third === BOM[2]) {
this.discardLeadingBytes(3)
}
this.checkBOM = false
return !this.hasCurrentByte()
}
} }
module.exports = { module.exports = {
@@ -38356,7 +38617,7 @@ function establishWebSocketConnection (url, protocols, client, ws, onEstablish,
// is specified, the server needs to include the same field and one of // is specified, the server needs to include the same field and one of
// the selected subprotocol values in its response for the connection to // the selected subprotocol values in its response for the connection to
// be established. // be established.
if (!requestProtocols.includes(secProtocol)) { if (requestProtocols === null || !requestProtocols.includes(secProtocol)) {
failWebsocketConnection(ws, 'Protocol was not set in the opening handshake.') failWebsocketConnection(ws, 'Protocol was not set in the opening handshake.')
return return
} }
@@ -39117,7 +39378,12 @@ class PerMessageDeflate {
if (this.#maxPayloadSize > 0 && this.#inflate[kLength] > this.#maxPayloadSize) { if (this.#maxPayloadSize > 0 && this.#inflate[kLength] > this.#maxPayloadSize) {
callback(new MessageSizeExceededError()) callback(new MessageSizeExceededError())
// The inflater may still hold buffered input that can emit a late
// zlib error. Remove the data listener, then deterministically stop
// the stream so a subsequent 'error' cannot fire without a listener
// (which would terminate the process as an unhandled error event).
this.#inflate.removeAllListeners() this.#inflate.removeAllListeners()
this.#inflate.destroy()
this.#inflate = null this.#inflate = null
return return
} }