From 0e0ab085a3c6338525fd69e4d22643936d6f0a9f Mon Sep 17 00:00:00 2001 From: Dawid Dziurla Date: Tue, 29 Sep 2026 16:54:26 +0200 Subject: [PATCH] node_modules: update (#328) Co-authored-by: dawidd6 <9713907+dawidd6@users.noreply.github.com> --- node_modules/.package-lock.json | 6 +- .../docs/docs/api/DiagnosticsChannel.md | 4 +- node_modules/undici/lib/core/request.js | 70 +++- .../undici/lib/dispatcher/client-h1.js | 41 +- .../undici/lib/dispatcher/client-h2.js | 90 +++- .../undici/lib/handler/retry-handler.js | 37 +- .../lib/web/eventsource/eventsource-stream.js | 395 +++++++++++------- .../undici/lib/web/websocket/connection.js | 2 +- .../lib/web/websocket/permessage-deflate.js | 5 + node_modules/undici/package.json | 2 +- 10 files changed, 460 insertions(+), 192 deletions(-) diff --git a/node_modules/.package-lock.json b/node_modules/.package-lock.json index 49b24142..750044d8 100644 --- a/node_modules/.package-lock.json +++ b/node_modules/.package-lock.json @@ -128,9 +128,9 @@ } }, "node_modules/undici": { - "version": "6.28.0", - "resolved": "https://registry.npmjs.org/undici/-/undici-6.28.0.tgz", - "integrity": "sha512-LIY910g9TI13YS95lrMFrs8Rm/u/irgHeTWoKCoteeJ04CUJ92eEfj0rVn+7VKMPBpUPiUoBKfhNyLI23EE/KA==", + "version": "6.29.0", + "resolved": "https://registry.npmjs.org/undici/-/undici-6.29.0.tgz", + "integrity": "sha512-R+RODBqp6i2pPflGdq+xIOUkl+RNfGgHwoinecKu/JCuf2uO06cOKoDbI2P7Dn6KcswdKwrczbU6IYJ6K8X+wg==", "license": "MIT", "engines": { "node": ">=18.17" diff --git a/node_modules/undici/docs/docs/api/DiagnosticsChannel.md b/node_modules/undici/docs/docs/api/DiagnosticsChannel.md index 099c072f..aefb1009 100644 --- a/node_modules/undici/docs/docs/api/DiagnosticsChannel.md +++ b/node_modules/undici/docs/docs/api/DiagnosticsChannel.md @@ -40,7 +40,8 @@ diagnosticsChannel.channel('undici:request:bodySent').subscribe(({ request }) => ## `undici:request:headers` -This message is published after the response headers have been received, i.e. the response has been completed. +This message is published after the response headers have been received. This includes a successful CONNECT or +protocol upgrade response. ```js import diagnosticsChannel from 'diagnostics_channel' @@ -57,6 +58,7 @@ diagnosticsChannel.channel('undici:request:headers').subscribe(({ request, respo ## `undici:request:trailers` This message is published after the response body and trailers have been received, i.e. the response has been completed. +After an upgraded socket is passed to the request handler, this message is published with an empty `trailers` array. ```js import diagnosticsChannel from 'diagnostics_channel' diff --git a/node_modules/undici/lib/core/request.js b/node_modules/undici/lib/core/request.js index 8e7ecc73..0b2b1a8d 100644 --- a/node_modules/undici/lib/core/request.js +++ b/node_modules/undici/lib/core/request.js @@ -263,11 +263,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.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) { diff --git a/node_modules/undici/lib/dispatcher/client-h1.js b/node_modules/undici/lib/dispatcher/client-h1.js index a801ecb3..52fd8d05 100644 --- a/node_modules/undici/lib/dispatcher/client-h1.js +++ b/node_modules/undici/lib/dispatcher/client-h1.js @@ -432,7 +432,7 @@ class Parser { } onUpgrade (head) { - const { upgrade, client, socket, headers, statusCode } = this + const { upgrade, client, socket, headers, statusCode, statusText } = this assert(upgrade) assert(client[kSocket] === socket) @@ -467,9 +467,10 @@ class Parser { client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade')) try { - request.onUpgrade(statusCode, headers, socket) - } catch (err) { - util.destroy(socket, err) + request.onUpgrade(statusCode, headers, socket, statusText) + } catch (error) { + util.errorRequest(client, request, error) + util.destroy(socket, error) } client[kResume]() @@ -876,7 +877,7 @@ async function connectH1 (client, socket) { function clearIdleSocketValidation (socket) { if (socket[kIdleSocketValidationTimeout]) { - clearTimeout(socket[kIdleSocketValidationTimeout]) + clearImmediate(socket[kIdleSocketValidationTimeout]) socket[kIdleSocketValidationTimeout] = null } @@ -885,15 +886,23 @@ function clearIdleSocketValidation (socket) { function scheduleIdleSocketValidation (client, socket) { 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[kIdleSocketValidation] = 2 if (client[kSocket] === socket && !socket.destroyed) { client[kResume]() } - }, 0) - socket[kIdleSocketValidationTimeout].unref?.() + }) } /** @@ -1042,12 +1051,22 @@ function writeH1 (client, request) { const socket = client[kSocket] clearIdleSocketValidation(socket) - const abort = (err) => { - if (request.aborted || request.completed) { + /** + * @param {Error} [error] + */ + const abort = (error) => { + if (request.aborted) { 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(socket, new InformationalError('aborted')) diff --git a/node_modules/undici/lib/dispatcher/client-h2.js b/node_modules/undici/lib/dispatcher/client-h2.js index 4a52effb..bb8eda84 100644 --- a/node_modules/undici/lib/dispatcher/client-h2.js +++ b/node_modules/undici/lib/dispatcher/client-h2.js @@ -1,6 +1,7 @@ 'use strict' const assert = require('node:assert') +const { errorMonitor } = require('node:events') const { pipeline } = require('node:stream') const util = require('../core/util.js') const { @@ -77,6 +78,15 @@ function parseH2Headers (headers) { 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) { client[kSocket] = socket @@ -297,22 +307,32 @@ function writeH2 (client, request) { headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}` headers[HTTP2_HEADER_METHOD] = method - const abort = (err) => { - if (request.aborted || request.completed) { + /** + * @param {Error} [error] + */ + const abort = (error) => { + if (request.aborted) { 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) { - util.destroy(stream, err) + util.destroy(stream, error) } // 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 - util.destroy(body, err) + util.destroy(body, error) client[kQueue][client[kRunningIdx]++] = null client[kResume]() } @@ -331,25 +351,57 @@ function writeH2 (client, request) { if (method === 'CONNECT') { 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 stream = session.request(headers, { endStream: false, signal }) + let upgradeResponseFinished = false - if (stream.id && !stream.pending) { - request.onUpgrade(null, null, stream) - ++session[kOpenStreams] - client[kQueue][client[kRunningIdx]++] = null - } else { - stream.once('ready', () => { - request.onUpgrade(null, null, stream) - ++session[kOpenStreams] - client[kQueue][client[kRunningIdx]++] = null - }) + /** + * @param {import('node:http2').IncomingHttpHeaders} headers + */ + const onResponse = (headers) => { + upgradeResponseFinished = true + stream.off(errorMonitor, onUpgradeError) + request.onUpgradeResponse(Number(headers[HTTP2_HEADER_STATUS]), headers, parseH2ResponseHeaders) } + /** + * @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', () => { + 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 if (session[kOpenStreams] === 0) session.unref() }) diff --git a/node_modules/undici/lib/handler/retry-handler.js b/node_modules/undici/lib/handler/retry-handler.js index c62a2409..7084c2b5 100644 --- a/node_modules/undici/lib/handler/retry-handler.js +++ b/node_modules/undici/lib/handler/retry-handler.js @@ -90,6 +90,7 @@ class RetryHandler { this.end = null this.etag = null this.resume = null + this.headersSent = false // Handle possible onConnect duplication this.handler.onConnect(reason => { @@ -102,6 +103,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 () { if (this.handler.onRequestSent) { this.handler.onRequestSent() @@ -190,7 +205,12 @@ class RetryHandler { this.retryCount += 1 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( statusCode, rawHeaders, @@ -259,8 +279,15 @@ class RetryHandler { const { start, size, end = size - 1 } = contentRange - assert(this.start === start, 'content-range mismatch') - assert(this.end == null || this.end === end, 'content-range mismatch') + if (this.start !== start || (this.end != null && this.end !== end)) { + this.abort( + new RequestRetryError('Content-Range mismatch', statusCode, { + headers, + data: { count: this.retryCount } + }) + ) + return false + } this.resume = resume return true @@ -272,6 +299,7 @@ class RetryHandler { const range = parseRangeHeader(headers['content-range']) if (range == null) { + this.headersSent = true return this.handler.onHeaders( statusCode, rawHeaders, @@ -310,6 +338,7 @@ class RetryHandler { ) this.resume = resume + this.headersSent = true this.etag = headers.etag != null ? headers.etag : null // Weak etags are not useful for comparison nor cache @@ -349,7 +378,7 @@ class RetryHandler { } 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) } diff --git a/node_modules/undici/lib/web/eventsource/eventsource-stream.js b/node_modules/undici/lib/web/eventsource/eventsource-stream.js index 75493456..af04e163 100644 --- a/node_modules/undici/lib/web/eventsource/eventsource-stream.js +++ b/node_modules/undici/lib/web/eventsource/eventsource-stream.js @@ -23,6 +23,49 @@ const COLON = 0x3A */ 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 * @type {object} @@ -63,11 +106,14 @@ class EventSourceStream extends Transform { eventEndCheck = false /** - * @type {Buffer} + * @type {Buffer[]} */ - buffer = null + chunks = [] + chunkIndex = 0 pos = 0 + lineChunkIndex = 0 + linePos = 0 event = { data: undefined, @@ -106,92 +152,20 @@ class EventSourceStream extends Transform { return } - // Cache the chunk in the buffer, as the data might not be complete while - // 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 - } + this.chunks.push(chunk) // Strip leading byte-order-mark if we opened the stream and started // the processing of the incoming data if (this.checkBOM) { - switch (this.buffer.length) { - case 1: - // Check if the first byte is the same as the first byte of the BOM - 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 + if (this.handleBOM()) { + callback() + return } } - 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 next character is also an end-of-line. if (this.eventEndCheck) { @@ -204,10 +178,9 @@ class EventSourceStream extends Transform { if (this.crlfCheck) { // If the current character is a line feed, we can remove it // from the buffer and reset the crlfCheck flag - if (this.buffer[this.pos] === LF) { - this.buffer = this.buffer.subarray(this.pos + 1) - this.pos = 0 + if (byte === LF) { this.crlfCheck = false + this.consumeCurrentByte() // It is possible that the line feed is not the end of the // event. We need to check if the next character is an @@ -223,19 +196,17 @@ class EventSourceStream extends Transform { 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 // 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 // buffer - if (this.buffer[this.pos] === CR) { + if (byte === CR) { this.crlfCheck = true } - this.buffer = this.buffer.subarray(this.pos + 1) - this.pos = 0 - if ( - this.event.data !== undefined || this.event.event || this.event.id || this.event.retry) { + this.consumeCurrentByte() + if (this.hasPendingEvent()) { this.processEvent(this.event) } this.clearEvent() @@ -249,22 +220,18 @@ class EventSourceStream extends Transform { // If the current character is an end-of-line, we can process the // 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 // set the crlfCheck flag to true, as we need to check if the // next character is a line feed - if (this.buffer[this.pos] === CR) { + if (byte === CR) { this.crlfCheck = true } // In any case, we can process the line as we reached an // end-of-line character - this.parseLine(this.buffer.subarray(0, this.pos), this.event) - - // 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 + this.parseLine(this.readLine(), this.event) + this.consumeCurrentByte() // 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 // finished. @@ -272,7 +239,7 @@ class EventSourceStream extends Transform { continue } - this.pos++ + this.advanceCursor() } callback() @@ -297,64 +264,53 @@ class EventSourceStream extends Transform { return } - let field = '' - let value = '' + let fieldLength = line.length + let valueStart = line.length // If the line contains a U+003A COLON character (:) if (colonPosition !== -1) { - // Collect the characters on the line before the first U+003A COLON - // 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') + fieldLength = colonPosition // Collect the characters on the line after the first U+003A COLON // character (:), and let value be that string. // If value starts with a U+0020 SPACE character, remove it from value. - let valueStart = colonPosition + 1 + valueStart = colonPosition + 1 if (line[valueStart] === SPACE) { ++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 - // decoded as UTF-8 - switch (field) { - case 'data': - if (event[field] === undefined) { - event[field] = value - } else { - event[field] += `\n${value}` - } - break - case 'retry': - if (isASCIINumber(value)) { - event[field] = value - } - break - case 'id': - if (isValidLastEventId(value)) { - event[field] = value - } - break - case 'event': - if (value.length > 0) { - event[field] = value - } - break + if (isFieldName(line, fieldLength, DATA)) { + const value = line.toString('utf8', valueStart) + + if (event.data === undefined) { + event.data = value + } else { + event.data += `\n${value}` + } + return + } + + if (isFieldName(line, fieldLength, RETRY)) { + if (isASCIINumberBytes(line, valueStart)) { + event.retry = line.toString('utf8', valueStart) + } + return + } + + if (isFieldName(line, fieldLength, ID)) { + if (isValidLastEventIdBytes(line, valueStart)) { + event.id = line.toString('utf8', valueStart) + } + return + } + + if (isFieldName(line, fieldLength, EVENT)) { + const value = line.toString('utf8', valueStart) + + if (value.length > 0) { + event.event = value + } } } @@ -384,13 +340,152 @@ class EventSourceStream extends Transform { } clearEvent () { - this.event = { - data: undefined, - event: undefined, - id: undefined, - retry: undefined + this.event.data = undefined + this.event.event = undefined + this.event.id = undefined + this.event.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 = { diff --git a/node_modules/undici/lib/web/websocket/connection.js b/node_modules/undici/lib/web/websocket/connection.js index bb87d361..1197b9b6 100644 --- a/node_modules/undici/lib/web/websocket/connection.js +++ b/node_modules/undici/lib/web/websocket/connection.js @@ -192,7 +192,7 @@ function establishWebSocketConnection (url, protocols, client, ws, onEstablish, // is specified, the server needs to include the same field and one of // the selected subprotocol values in its response for the connection to // be established. - if (!requestProtocols.includes(secProtocol)) { + if (requestProtocols === null || !requestProtocols.includes(secProtocol)) { failWebsocketConnection(ws, 'Protocol was not set in the opening handshake.') return } diff --git a/node_modules/undici/lib/web/websocket/permessage-deflate.js b/node_modules/undici/lib/web/websocket/permessage-deflate.js index 6a6e4389..0b3d493d 100644 --- a/node_modules/undici/lib/web/websocket/permessage-deflate.js +++ b/node_modules/undici/lib/web/websocket/permessage-deflate.js @@ -63,7 +63,12 @@ class PerMessageDeflate { if (this.#maxPayloadSize > 0 && this.#inflate[kLength] > this.#maxPayloadSize) { 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.destroy() this.#inflate = null return } diff --git a/node_modules/undici/package.json b/node_modules/undici/package.json index cd46deca..1b87e5eb 100644 --- a/node_modules/undici/package.json +++ b/node_modules/undici/package.json @@ -1,6 +1,6 @@ { "name": "undici", - "version": "6.28.0", + "version": "6.29.0", "description": "An HTTP/1.1 client, written from scratch for Node.js", "homepage": "https://undici.nodejs.org", "bugs": {