From 084bc5f315e40ca7a7c23a480c0a6dfa1d1ace54 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:29:39 +0200 Subject: [PATCH] build(deps): bump undici from 7.29.0 to 7.30.0 (#1425) Co-authored-by: Fernandez Ludovic --- dist/post_run/index.js | 1225 ++++++++++++++++++++++++++++++---------- dist/run/index.js | 1225 ++++++++++++++++++++++++++++++---------- package-lock.json | 6 +- 3 files changed, 1843 insertions(+), 613 deletions(-) diff --git a/dist/post_run/index.js b/dist/post_run/index.js index 78f8856..9e36060 100644 --- a/dist/post_run/index.js +++ b/dist/post_run/index.js @@ -34858,11 +34858,32 @@ class Request { } } - onUpgrade (statusCode, headers, socket) { + /** + * @param {number} statusCode + * @param {Buffer[]|string[]} 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 (channels.headers.hasSubscribers) { + channels.headers.publish({ request: this, response: { statusCode, headers, statusText } }) + } + + const result = this[kHandler].onUpgrade(statusCode, headers, socket) + + if (!this.aborted) { + this.completed = true + if (channels.trailers.hasSubscribers) { + channels.trailers.publish({ request: this, trailers: [] }) + } + } + + return result } onComplete (trailers) { @@ -37144,14 +37165,16 @@ function defaultFactory (origin, opts) { } class BalancedPool extends PoolBase { - constructor (upstreams = [], { factory = defaultFactory, ...opts } = {}) { + constructor (upstreams = [], { factory = defaultFactory, connect, tls, ...opts } = {}) { if (typeof factory !== 'function') { throw new InvalidArgumentError('factory must be a function.') } super(opts) - this[kOptions] = { ...util.deepClone(opts) } + if (connect && typeof connect !== 'function') connect = { ...connect } + if (tls && typeof tls !== 'function') tls = { ...tls } + this[kOptions] = { ...util.deepClone(opts), connect, tls } this[kOptions].interceptors = opts.interceptors ? { ...opts.interceptors } : undefined @@ -37391,8 +37414,14 @@ function lazyllhttp () { let mod - // We disable wasm SIMD on ppc64 as it seems to be broken on Power 9 architectures. - let useWasmSIMD = process.arch !== 'ppc64' + // We disable wasm SIMD on older versions of Node.js on ppc64 that are broken on Power >=9 architectures. + let useWasmSIMD = true + if (process.arch === 'ppc64') { + const [major, minor] = process.versions.node.split('.').map(n => parseInt(n, 10)) + if (major < 24 || (major === 24 && minor < 12)) { + useWasmSIMD = false + } + } // The Env Variable UNDICI_NO_WASM_SIMD allows explicitly overriding the default behavior if (process.env.UNDICI_NO_WASM_SIMD === '1') { useWasmSIMD = false @@ -37843,7 +37872,7 @@ class Parser { * @param {Buffer} head */ onUpgrade (head) { - const { upgrade, client, socket, headers, statusCode } = this + const { upgrade, client, socket, headers, statusCode, statusText } = this assert(upgrade) assert(client[kSocket] === socket) @@ -37878,8 +37907,9 @@ class Parser { client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade')) try { - request.onUpgrade(statusCode, headers, socket) + request.onUpgrade(statusCode, headers, socket, statusText) } catch (err) { + util.errorRequest(client, request, err) util.destroy(socket, err) } @@ -38334,7 +38364,7 @@ function onSocketClose () { function clearIdleSocketValidation (socket) { if (socket[kIdleSocketValidationTimeout]) { - clearTimeout(socket[kIdleSocketValidationTimeout]) + clearImmediate(socket[kIdleSocketValidationTimeout]) socket[kIdleSocketValidationTimeout] = null } @@ -38343,15 +38373,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?.() + }) } /** @@ -39069,7 +39107,9 @@ const { RequestAbortedError, SocketError, InformationalError, - InvalidArgumentError + InvalidArgumentError, + HeadersTimeoutError, + BodyTimeoutError } = __nccwpck_require__(68707) const { kUrl, @@ -39094,6 +39134,7 @@ const { kHTTPContext, kClosed, kBodyTimeout, + kHeadersTimeout, kEnableConnectProtocol, kRemoteSettings, kHTTP2Stream, @@ -39280,7 +39321,11 @@ function resumeH2 (client) { const socket = client[kSocket] if (socket?.destroyed === false) { - if (client[kSize] === 0 || client[kMaxConcurrentStreams] === 0) { + // Only let the process exit when there is genuinely nothing outstanding. + // Unreffing because the peer advertised MAX_CONCURRENT_STREAMS = 0 left + // queued requests with nothing holding the event loop open, so the process + // could exit with status 0 while an awaited request never settled. + if (client[kSize] === 0) { socket.unref() client[kHTTP2Session].unref() } else { @@ -39375,6 +39420,36 @@ function onHttp2SessionEnd () { * @this {import('http2').ClientHttp2Session} * @param {number} errorCode */ +// Backport of #5410 and #5569. HTTP/2 multiplexes, so requests complete out of +// order; advancing kRunningIdx blindly retired whichever request happened to +// sit at the head instead of the one that actually finished, which both lost +// requests and left phantom running slots behind. +function completeRequest (client, request, resetPendingIdx = false) { + const queue = client[kQueue] + const runningIdx = client[kRunningIdx] + + // In-order completion: clear the request and advance without splicing. + // The client's resume loop compacts cleared slots once the index grows. + if (runningIdx < client[kPendingIdx] && queue[runningIdx] === request) { + queue[runningIdx] = null + client[kRunningIdx] = runningIdx + 1 + return + } + + const index = queue.indexOf(request, runningIdx) + + if (index === -1 || index >= client[kPendingIdx]) { + return + } + + queue.splice(index, 1) + client[kPendingIdx]-- + + if (resetPendingIdx && client[kPendingIdx] < client[kRunningIdx]) { + client[kPendingIdx] = client[kRunningIdx] + } +} + function onHttp2SessionGoAway (errorCode) { // TODO(mcollina): Verify if GOAWAY implements the spec correctly: // https://datatracker.ietf.org/doc/html/rfc7540#section-6.8 @@ -39396,7 +39471,9 @@ function onHttp2SessionGoAway (errorCode) { if (client[kRunningIdx] < client[kQueue].length) { const request = client[kQueue][client[kRunningIdx]] client[kQueue][client[kRunningIdx]++] = null - util.errorRequest(client, request, err) + if (request != null) { + util.errorRequest(client, request, err) + } client[kPendingIdx] = client[kRunningIdx] } @@ -39429,7 +39506,9 @@ function onHttp2SessionClose () { const requests = client[kQueue].splice(client[kRunningIdx]) for (let i = 0; i < requests.length; i++) { const request = requests[i] - util.errorRequest(client, request, err) + if (request != null) { + util.errorRequest(client, request, err) + } } } } @@ -39477,7 +39556,10 @@ function shouldSendContentLength (method) { } function writeH2 (client, request) { - const requestTimeout = request.bodyTimeout ?? client[kBodyTimeout] + // Time to the response headers, then time between body chunks. Using + // bodyTimeout for both made headersTimeout a no-op over HTTP/2. + const headersTimeout = request.headersTimeout ?? client[kHeadersTimeout] + const bodyTimeout = request.bodyTimeout ?? client[kBodyTimeout] const session = client[kHTTP2Session] const { method, path, host, upgrade, expectContinue, signal, protocol, headers: reqHeaders } = request let { body } = request @@ -39544,6 +39626,7 @@ function writeH2 (client, request) { // We move the running index to the next request client[kOnError](err) + completeRequest(client, request) client[kResume]() } @@ -39591,14 +39674,18 @@ function writeH2 (client, request) { stream = session.request(headers, { endStream: false, signal }) stream[kHTTP2Stream] = true + ++session[kOpenStreams] stream.once('response', (headers, _flags) => { const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers - request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) - - ++session[kOpenStreams] - client[kQueue][client[kRunningIdx]++] = null + try { + request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) + } catch (err) { + abort(err) + return + } + completeRequest(client, request) }) stream.on('error', () => { @@ -39615,7 +39702,7 @@ function writeH2 (client, request) { if (session[kOpenStreams] === 0) session.unref() }) - stream.setTimeout(requestTimeout) + stream.setTimeout(headersTimeout) return true } @@ -39626,18 +39713,24 @@ function writeH2 (client, request) { // We disabled endStream to allow the user to write to the stream stream = session.request(headers, { endStream: false, signal }) stream[kHTTP2Stream] = true + ++session[kOpenStreams] stream.on('response', headers => { const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers - request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) - ++session[kOpenStreams] - client[kQueue][client[kRunningIdx]++] = null + try { + request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) + } catch (err) { + abort(err) + return + } + completeRequest(client, request) }) + stream.on('error', abort) stream.once('close', () => { session[kOpenStreams] -= 1 if (session[kOpenStreams] === 0) session.unref() }) - stream.setTimeout(requestTimeout) + stream.setTimeout(headersTimeout) return true } @@ -39738,7 +39831,7 @@ function writeH2 (client, request) { // Increment counter as we have new streams open ++session[kOpenStreams] - stream.setTimeout(requestTimeout) + stream.setTimeout(headersTimeout) // Track whether we received a response (headers) let responseReceived = false @@ -39747,6 +39840,7 @@ function writeH2 (client, request) { const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers request.onResponseStarted() responseReceived = true + stream.setTimeout(bodyTimeout) // Due to the stream nature, it is possible we face a race condition // where the stream has been assigned, but the request has been aborted @@ -39781,14 +39875,13 @@ function writeH2 (client, request) { request.onComplete({}) } - client[kQueue][client[kRunningIdx]++] = null + completeRequest(client, request) client[kResume]() } else { // Stream ended without receiving a response - this is an error // (e.g., server destroyed the stream before sending headers) abort(new InformationalError('HTTP/2: stream half-closed (remote)')) - client[kQueue][client[kRunningIdx]++] = null - client[kPendingIdx] = client[kRunningIdx] + completeRequest(client, request, true) client[kResume]() } }) @@ -39799,6 +39892,14 @@ function writeH2 (client, request) { if (session[kOpenStreams] === 0) { session.unref() } + + // A stream can close without ever emitting 'end' or 'error': a peer's + // RST_STREAM(CANCEL) received before the response is reported by Node as a + // bare 'close', and destroying the stream unenrolls its timeout, so no + // 'timeout' follows either. Nothing else would ever settle this request. + if (!request.aborted && !request.completed) { + abort(new InformationalError('HTTP/2: stream closed before the response was complete')) + } }) stream.once('error', function (err) { @@ -39816,7 +39917,9 @@ function writeH2 (client, request) { }) stream.on('timeout', () => { - const err = new InformationalError(`HTTP/2: "stream timeout after ${requestTimeout}"`) + const err = responseReceived + ? new BodyTimeoutError(`HTTP/2: "body timeout after ${bodyTimeout}"`) + : new HeadersTimeoutError(`HTTP/2: "headers timeout after ${headersTimeout}"`) stream.removeAllListeners('data') session[kOpenStreams] -= 1 @@ -40438,7 +40541,9 @@ class Client extends DispatcherBase { const requests = this[kQueue].splice(this[kPendingIdx]) for (let i = 0; i < requests.length; i++) { const request = requests[i] - util.errorRequest(this, request, err) + if (request != null) { + util.errorRequest(this, request, err) + } } const callback = () => { @@ -40477,7 +40582,9 @@ function onError (client, err) { for (let i = 0; i < requests.length; i++) { const request = requests[i] - util.errorRequest(client, request, err) + if (request != null) { + util.errorRequest(client, request, err) + } } assert(client[kSize] === 0) } @@ -42786,6 +42893,13 @@ class CacheHandler { } const cacheControlHeader = resHeaders['cache-control'] + const cacheControlDirectives = cacheControlHeader ? parseCacheControlHeader(cacheControlHeader) : {} + + if (revalidationResponseDisallowsCachedReuse(this.#cacheType, resHeaders, cacheControlDirectives)) { + deleteCachedValue(this.#store, this.#cacheKey) + return downstreamOnHeaders() + } + const heuristicallyCacheable = resHeaders['last-modified'] && arrayIncludes(HEURISTICALLY_CACHEABLE_STATUS_CODES, statusCode) if ( !cacheControlHeader && @@ -42802,8 +42916,7 @@ class CacheHandler { return downstreamOnHeaders() } - const cacheControlDirectives = cacheControlHeader ? parseCacheControlHeader(cacheControlHeader) : {} - if (!canCacheResponse(this.#cacheType, statusCode, resHeaders, cacheControlDirectives, this.#cacheKey.headers)) { + if (!canCacheResponse(this.#cacheType, this.#cacheKey.method, statusCode, resHeaders, cacheControlDirectives, this.#cacheKey.headers)) { if (statusCode === 304 && (cacheControlHeader || revalidationResponseDisallowsCachedReuse(this.#cacheType, resHeaders, cacheControlDirectives))) { deleteCachedValue(this.#store, this.#cacheKey) } @@ -43044,7 +43157,10 @@ function deleteCachedValueIfNotModified (statusCode, store, cacheKey) { */ function revalidationResponseDisallowsCachedReuse (cacheType, resHeaders, cacheControlDirectives) { return cacheControlDirectives['no-store'] === true || - (cacheType === 'shared' && cacheControlDirectives.private === true) || + (cacheType === 'shared' && ( + cacheControlDirectives.private === true || + Object.hasOwn(resHeaders, 'set-cookie') + )) || (resHeaders.vary ? isInvalidOrWildcardVaryHeader(resHeaders.vary) : false) } @@ -43052,12 +43168,16 @@ function revalidationResponseDisallowsCachedReuse (cacheType, resHeaders, cacheC * @see https://www.rfc-editor.org/rfc/rfc9111.html#name-storing-responses-to-authen * * @param {import('../../types/cache-interceptor.d.ts').default.CacheOptions['type']} cacheType + * @param {string} method * @param {number} statusCode * @param {import('../../types/header.d.ts').IncomingHttpHeaders} resHeaders * @param {import('../../types/cache-interceptor.d.ts').default.CacheControlDirectives} cacheControlDirectives * @param {import('../../types/header.d.ts').IncomingHttpHeaders} [reqHeaders] */ -function canCacheResponse (cacheType, statusCode, resHeaders, cacheControlDirectives, reqHeaders) { +function canCacheResponse (cacheType, method, statusCode, resHeaders, cacheControlDirectives, reqHeaders) { + if (!arrayIncludes(util.safeHTTPMethods, method)) { + return false + } // Status code must be final and understood. if (statusCode < 200 || arrayIncludes(NOT_UNDERSTOOD_STATUS_CODES, statusCode)) { return false @@ -43078,7 +43198,10 @@ function canCacheResponse (cacheType, statusCode, resHeaders, cacheControlDirect return false } - if (cacheType === 'shared' && cacheControlDirectives.private === true) { + if (cacheType === 'shared' && ( + cacheControlDirectives.private === true || + Object.hasOwn(resHeaders, 'set-cookie') + )) { return false } @@ -44303,6 +44426,55 @@ function validatePartialResponseContentLength (headers, range, statusCode, retry } } +// A stable controller handed to the downstream handler for the lifetime of the +// request. Each transparent retry/resume is a separate dispatch with its own +// connection controller. The proxy always forwards to the active connection +// while preserving a downstream pause across controller replacement. +class RetryController { + #paused = false + #target = null + + set target (target) { + this.#target = target + if (this.#paused) { + target?.pause() + } + } + + get target () { return this.#target } + + pause () { + this.#paused = true + this.#target?.pause() + } + + resume () { + this.#paused = false + this.#target?.resume() + } + + abort (reason) { + this.#target?.abort(reason) + } + + get paused () { return this.#paused || (this.#target?.paused ?? false) } + get aborted () { return this.#target?.aborted ?? false } + get reason () { return this.#target?.reason ?? null } + get rawHeaders () { return this.#target?.rawHeaders ?? null } + set rawHeaders (value) { + if (this.#target) { + this.#target.rawHeaders = value + } + } + + get rawTrailers () { return this.#target?.rawTrailers ?? null } + set rawTrailers (value) { + if (this.#target) { + this.#target.rawTrailers = value + } + } +} + class RetryHandler { constructor (opts, { dispatch, handler }) { const { retryOptions, ...dispatchOpts } = opts @@ -44357,14 +44529,23 @@ class RetryHandler { this.start = 0 this.end = null this.etag = null + this.controllerProxy = new RetryController() } onResponseStartWithRetry (controller, statusCode, headers, statusMessage, err) { if (this.retryOpts.throwOnError) { // Preserve old behavior for status codes that are not eligible for retry if (this.retryOpts.statusCodes.includes(statusCode) === false) { - this.headersSent = true - this.handler.onResponseStart?.(controller, statusCode, headers, statusMessage) + if (this.headersSent) { + // The downstream handler already received the response from an + // earlier attempt. Forwarding this response would replace the + // downstream body and leave the original body pending forever. + this.handler.onResponseError?.(this.controllerProxy, err) + } else { + this.headersSent = true + this.checkpointResponseEnd(headers) + this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage) + } } else { this.error = err } @@ -44374,14 +44555,23 @@ class RetryHandler { if (isDisturbed(this.opts.body)) { this.headersSent = true - this.handler.onResponseStart?.(controller, statusCode, headers, statusMessage) + this.checkpointResponseEnd(headers) + this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage) return } function shouldRetry (passedErr) { if (passedErr) { - this.headersSent = true - this.handler.onResponseStart?.(controller, statusCode, headers, statusMessage) + if (this.headersSent) { + // The downstream handler already received the response from an + // earlier attempt. Forwarding this response would replace the + // downstream body and leave the original body pending forever. + this.handler.onResponseError?.(this.controllerProxy, passedErr) + } else { + this.headersSent = true + this.checkpointResponseEnd(headers) + this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage) + } controller.resume() return } @@ -44401,14 +44591,29 @@ class RetryHandler { ) } - onRequestStart (controller, context) { - if (!this.headersSent) { - this.handler.onRequestStart?.(controller, context) + checkpointResponseEnd (headers) { + 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 } } - onRequestUpgrade (controller, statusCode, headers, socket) { - this.handler.onRequestUpgrade?.(controller, statusCode, headers, socket) + onRequestStart (controller, context) { + this.controllerProxy.target = controller + if (!this.headersSent) { + this.handler.onRequestStart?.(this.controllerProxy, context) + } + } + + onRequestUpgrade (_controller, statusCode, headers, socket) { + this.handler.onRequestUpgrade?.(this.controllerProxy, statusCode, headers, socket) } static [kRetryHandlerDefaultRetry] (err, { state, opts }, cb) { @@ -44521,8 +44726,12 @@ class RetryHandler { const { start, size, end = size ? size - 1 : null } = 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)) { + throw new RequestRetryError('Content-Range mismatch', statusCode, { + headers, + data: { count: this.retryCount } + }) + } return } @@ -44535,7 +44744,7 @@ class RetryHandler { if (range == null) { this.headersSent = true this.handler.onResponseStart?.( - controller, + this.controllerProxy, statusCode, headers, statusMessage @@ -44584,7 +44793,7 @@ class RetryHandler { this.headersSent = true this.handler.onResponseStart?.( - controller, + this.controllerProxy, statusCode, headers, statusMessage @@ -44597,30 +44806,30 @@ class RetryHandler { } } - onResponseData (controller, chunk) { + onResponseData (_controller, chunk) { if (this.error) { return } this.start += chunk.length - this.handler.onResponseData?.(controller, chunk) + this.handler.onResponseData?.(this.controllerProxy, chunk) } - onResponseEnd (controller, trailers) { + onResponseEnd (_controller, trailers) { if (this.error && this.retryOpts.throwOnError) { throw this.error } if (!this.error) { this.retryCount = 0 - return this.handler.onResponseEnd?.(controller, trailers) + return this.handler.onResponseEnd?.(this.controllerProxy, trailers) } - this.retry(controller) + this.retry() } - retry (controller) { + retry () { if (this.start !== 0) { const headers = { range: `bytes=${this.start}-${this.end ?? ''}` } @@ -44642,23 +44851,23 @@ class RetryHandler { this.retryCountCheckpoint = this.retryCount this.dispatch(this.opts, this) } catch (err) { - this.handler.onResponseError?.(controller, err) + this.handler.onResponseError?.(this.controllerProxy, err) } } onResponseError (controller, err) { - if (controller?.aborted || isDisturbed(this.opts.body)) { - this.handler.onResponseError?.(controller, err) + if (controller?.aborted || isDisturbed(this.opts.body) || (this.headersSent && !this.resume)) { + this.handler.onResponseError?.(this.controllerProxy, err) return } function shouldRetry (returnedErr) { if (!returnedErr) { - this.retry(controller) + this.retry() return } - this.handler?.onResponseError?.(controller, returnedErr) + this.handler?.onResponseError?.(this.controllerProxy, returnedErr) } // We reconcile in case of a mix between network errors @@ -45038,7 +45247,10 @@ function staleResponseRequiresRevalidation (result, cacheType) { * @returns {boolean} */ function revalidationResponseDisallowsCachedReuse (cacheType, headers) { - if (headers.vary && isInvalidOrWildcardVaryHeader(headers.vary)) { + if ( + (headers.vary && isInvalidOrWildcardVaryHeader(headers.vary)) || + (cacheType === 'shared' && Object.hasOwn(headers, 'set-cookie')) + ) { return true } @@ -45297,6 +45509,17 @@ function handleResult ( return handleUncachedResponse(dispatch, globalOpts, cacheKey, handler, opts, reqCacheControl) } + // Shared stores may outlive the Undici version that wrote them. Do not + // re-serve a Set-Cookie header from an existing shared-cache entry. + if (globalOpts.type === 'shared' && Object.hasOwn(result.headers, 'set-cookie')) { + if (util.isStream(result.body)) { + result.body.on('error', nop).destroy() + } + + deleteCachedValue(globalOpts.store, cacheKey) + return handleUncachedResponse(dispatch, globalOpts, cacheKey, handler, opts, reqCacheControl) + } + const now = Date.now() if (now > result.deleteAt) { // Response is expired, cache store shouldn't have given this to us @@ -45495,6 +45718,11 @@ module.exports = (opts = {}) => { * @type {import('../../types/cache-interceptor.d.ts').default.CacheKey} */ const cacheKey = makeCacheKey(opts) + + if (!arrayIncludes(util.safeHTTPMethods, opts.method)) { + return dispatch(opts, new CacheHandler(globalOpts, cacheKey, handler)) + } + const result = store.get(cacheKey) if (result && typeof result.then === 'function') { @@ -45532,7 +45760,8 @@ module.exports = (opts = {}) => { const { createInflate, createGunzip, createBrotliDecompress, createZstdDecompress } = __nccwpck_require__(38522) -const { pipeline } = __nccwpck_require__(57075) +const { pipeline, Transform: TransformStream } = __nccwpck_require__(57075) +const { InvalidArgumentError, ResponseExceededMaxSizeError } = __nccwpck_require__(68707) const DecoratorHandler = __nccwpck_require__(58155) const { runtimeFeatures } = __nccwpck_require__(313) @@ -45540,6 +45769,60 @@ const { runtimeFeatures } = __nccwpck_require__(313) /** @typedef {import('node:stream').Transform} Controller */ /** @typedef {Transform&import('node:zlib').Zlib} DecompressorStream */ +class DecompressController { + #onPause + #onResume + #onAbort + #paused = false + + constructor (onPause, onResume, onAbort) { + this.#onPause = onPause + this.#onResume = onResume + this.#onAbort = onAbort + this.target = null + } + + pause () { + if (this.#paused) { + return + } + + this.#paused = true + this.#onPause() + } + + resume () { + if (!this.#paused) { + return + } + + this.#paused = false + this.#onResume() + } + + abort (reason) { + this.target?.abort(reason) + this.#onAbort(reason) + } + + get paused () { return this.#paused } + get aborted () { return this.target?.aborted ?? false } + get reason () { return this.target?.reason ?? null } + get rawHeaders () { return this.target?.rawHeaders ?? null } + set rawHeaders (value) { + if (this.target) { + this.target.rawHeaders = value + } + } + + get rawTrailers () { return this.target?.rawTrailers ?? null } + set rawTrailers (value) { + if (this.target) { + this.target.rawTrailers = value + } + } +} + /** @type {Record DecompressorStream>} */ const supportedEncodings = { gzip: createGunzip, @@ -45552,6 +45835,31 @@ const supportedEncodings = { } const defaultSkipStatusCodes = /** @type {const} */ ([204, 304]) +const defaultMaxSize = 0 + +/** + * Limits the output of one stage in a decompression chain. + * @param {number} maxSize - Maximum output size in bytes + * @returns {Transform} + */ +function createMaxSizeLimiter (maxSize) { + let size = 0 + + return new TransformStream({ + transform (chunk, _encoding, callback) { + const decompressedSize = size + chunk.length + if (decompressedSize > maxSize) { + callback(new ResponseExceededMaxSizeError( + `Decompressed response size (${decompressedSize}) exceeded maxSize (${maxSize})` + )) + return + } + + size = decompressedSize + callback(null, chunk) + } + }) +} let warningEmitted = /** @type {boolean} */ (false) @@ -45559,20 +45867,150 @@ let warningEmitted = /** @type {boolean} */ (false) * @typedef {Object} DecompressHandlerOptions * @property {number[]|Readonly} [skipStatusCodes=[204, 304]] - List of status codes to skip decompression for * @property {boolean} [skipErrorResponses] - Whether to skip decompression for error responses (status codes >= 400) + * @property {number} [maxSize=0] - Maximum decompressed response size in bytes. 0 disables the limit */ class DecompressHandler extends DecoratorHandler { /** @type {Transform[]} */ #decompressors = [] + /** @type {Record | undefined} */ + #trailers /** @type {Readonly} */ #skipStatusCodes /** @type {boolean} */ #skipErrorResponses + /** @type {number} */ + #maxSize + /** @type {number} */ + #decompressedSize = 0 + /** @type {boolean} */ + #terminated = false + /** @type {boolean} */ + #inputEnded = false + /** @type {boolean} */ + #inputBackpressured = false + /** @type {boolean} */ + #upstreamPaused = false + /** @type {boolean} */ + #draining = false + /** @type {boolean} */ + #drainRequested = false + /** @type {boolean} */ + #completionPending = false + /** @type {DecompressorStream | undefined} */ + #finalDecompressor + /** @type {DecompressController} */ + #controller + + constructor (handler, { skipStatusCodes = defaultSkipStatusCodes, skipErrorResponses = true, maxSize = defaultMaxSize } = {}) { + if (!Number.isSafeInteger(maxSize) || maxSize < 0) { + throw new InvalidArgumentError('maxSize must be a non-negative integer') + } - constructor (handler, { skipStatusCodes = defaultSkipStatusCodes, skipErrorResponses = true } = {}) { super(handler) this.#skipStatusCodes = skipStatusCodes this.#skipErrorResponses = skipErrorResponses + this.#maxSize = maxSize + this.#controller = new DecompressController( + () => this.#onDownstreamPause(), + () => this.#onDownstreamResume(), + reason => { + if (this.#inputEnded && !this.#terminated) { + this.onResponseError(this.#controller, reason) + } + } + ) + } + + #onDownstreamPause () { + this.#pauseUpstream() + } + + #onDownstreamResume () { + const drainWasDeferred = this.#draining + this.#drainOutput() + if (!drainWasDeferred) { + this.#resumeUpstreamIfNeeded() + this.#finishIfReady() + } + } + + #pauseUpstream () { + if (!this.#upstreamPaused && !this.#terminated) { + this.#upstreamPaused = true + this.#controller.target?.pause() + } + } + + #resumeUpstreamIfNeeded () { + if (this.#upstreamPaused && !this.#controller.paused && !this.#inputBackpressured) { + this.#upstreamPaused = false + if (!this.#inputEnded) { + this.#controller.target?.resume() + } + } + } + + #drainOutput () { + if (this.#terminated || this.#controller.paused || !this.#finalDecompressor) { + return + } + + if (this.#draining) { + this.#drainRequested = true + return + } + + this.#draining = true + try { + do { + this.#drainRequested = false + let chunk + while (!this.#terminated && !this.#controller.paused && (chunk = this.#finalDecompressor.read()) !== null) { + if (this.#maxSize > 0) { + const decompressedSize = this.#decompressedSize + chunk.length + if (decompressedSize > this.#maxSize) { + this.#fail(new ResponseExceededMaxSizeError( + `Decompressed response size (${decompressedSize}) exceeded maxSize (${this.#maxSize})` + )) + return + } + + this.#decompressedSize = decompressedSize + } + + const result = super.onResponseData(this.#controller, chunk) + if (result === false && !this.#controller.paused) { + this.#controller.pause() + } + } + } while (this.#drainRequested && !this.#terminated && !this.#controller.paused) + } finally { + this.#draining = false + } + + this.#resumeUpstreamIfNeeded() + this.#finishIfReady() + } + + #finishIfReady () { + if (this.#terminated || !this.#completionPending || this.#controller.paused || this.#draining) { + return + } + + this.#terminated = true + this.#cleanupDecompressors() + super.onResponseEnd(this.#controller, this.#trailers) + } + + #onDecompressionEnd () { + if (this.#terminated) { + return + } + + this.#completionPending = true + this.#drainOutput() + this.#finishIfReady() } /** @@ -45592,7 +46030,7 @@ class DecompressHandler extends DecoratorHandler { * Creates a chain of decompressors for multiple content encodings * * @param {string} encodings - Comma-separated list of content encodings - * @returns {Array} - Array of decompressor streams + * @returns {Array} - Array of decompressor and limiting streams * @throws {Error} - If the number of content-encodings exceeds the maximum allowed */ #createDecompressionChain (encodings) { @@ -45620,60 +46058,97 @@ class DecompressHandler extends DecoratorHandler { decompressors.push(supportedEncodings[encoding]()) } - return decompressors + if (decompressors.length < 2) { + return decompressors + } + + /** @type {Transform[]} */ + const streams = [] + for (let i = 0; i < decompressors.length; i++) { + streams.push(decompressors[i]) + if (i < decompressors.length - 1 && this.#maxSize > 0) { + streams.push(createMaxSizeLimiter(this.#maxSize)) + } + } + + return streams } /** - * Sets up event handlers for a decompressor stream using readable events - * @param {DecompressorStream} decompressor - The decompressor stream - * @param {Controller} controller - The controller to coordinate with + * Stops decompression and reports an error. + * @param {Error} error - The decompression error * @returns {void} */ - #setupDecompressorEvents (decompressor, controller) { - decompressor.on('readable', () => { - let chunk - while ((chunk = decompressor.read()) !== null) { - const result = super.onResponseData(controller, chunk) - if (result === false) { - break - } - } - }) + #fail (error) { + if (this.#terminated) { + return + } - decompressor.on('error', (error) => { - super.onResponseError(controller, error) - }) + if (this.#inputEnded) { + // The request is already marked complete once the compressed input ends, + // so controller.abort() can no longer propagate decoder flush errors. + this.onResponseError(this.#controller, error) + } else { + this.#controller.abort(error) + } + } + + /** + * Sets up event handlers for the final decompressor stream. + * @param {DecompressorStream} decompressor - The decompressor stream + * @returns {void} + */ + #setupDecompressorEvents (decompressor) { + this.#finalDecompressor = decompressor + decompressor.on('readable', () => this.#drainOutput()) + decompressor.on('error', (error) => this.#fail(error)) } /** * Sets up event handling for a single decompressor - * @param {Controller} controller - The controller to handle events * @returns {void} */ - #setupSingleDecompressor (controller) { + #setupSingleDecompressor () { const decompressor = this.#decompressors[0] - this.#setupDecompressorEvents(decompressor, controller) + this.#setupDecompressorEvents(decompressor) - decompressor.on('end', () => { - super.onResponseEnd(controller, {}) - }) + decompressor.on('end', () => this.#onDecompressionEnd()) } /** * Sets up event handling for multiple chained decompressors using pipeline - * @param {Controller} controller - The controller to handle events * @returns {void} */ - #setupMultipleDecompressors (controller) { + #setupMultipleDecompressors () { const lastDecompressor = this.#decompressors[this.#decompressors.length - 1] - this.#setupDecompressorEvents(lastDecompressor, controller) + this.#setupDecompressorEvents(lastDecompressor) pipeline(this.#decompressors, (err) => { - if (err) { - super.onResponseError(controller, err) + if (this.#terminated) { return } - super.onResponseEnd(controller, {}) + + if (err) { + this.#fail(err) + return + } + + this.#onDecompressionEnd() + }) + } + + #setupInputBackpressure () { + const decompressor = this.#decompressors[0] + decompressor.on('drain', () => { + if (this.#terminated) { + return + } + + this.#inputBackpressured = false + if (!this.#controller.paused) { + this.#drainOutput() + this.#resumeUpstreamIfNeeded() + } }) } @@ -45683,6 +46158,16 @@ class DecompressHandler extends DecoratorHandler { */ #cleanupDecompressors () { this.#decompressors.length = 0 + this.#finalDecompressor = undefined + } + + onRequestStart (controller, context) { + this.#controller.target = controller + return super.onRequestStart(this.#controller, context) + } + + onRequestUpgrade (controller, statusCode, headers, socket) { + return super.onRequestUpgrade(this.#controller, statusCode, headers, socket) } /** @@ -45697,14 +46182,14 @@ class DecompressHandler extends DecoratorHandler { // If content encoding is not supported or status code is in skip list if (this.#shouldSkipDecompression(contentEncoding, statusCode)) { - return super.onResponseStart(controller, statusCode, headers, statusMessage) + return super.onResponseStart(this.#controller, statusCode, headers, statusMessage) } const decompressors = this.#createDecompressionChain(contentEncoding.toLowerCase()) if (decompressors.length === 0) { this.#cleanupDecompressors() - return super.onResponseStart(controller, statusCode, headers, statusMessage) + return super.onResponseStart(this.#controller, statusCode, headers, statusMessage) } this.#decompressors = decompressors @@ -45712,13 +46197,41 @@ class DecompressHandler extends DecoratorHandler { // Remove compression headers since we're decompressing const { 'content-encoding': _, 'content-length': __, ...newHeaders } = headers - if (this.#decompressors.length === 1) { - this.#setupSingleDecompressor(controller) - } else { - this.#setupMultipleDecompressors(controller) + if (this.#controller.rawHeaders) { + const rawHeaders = this.#controller.rawHeaders + + if (Array.isArray(rawHeaders)) { + const filteredHeaders = [] + for (let i = 0; i < rawHeaders.length; i += 2) { + const headerName = rawHeaders[i] + const name = Buffer.isBuffer(headerName) ? headerName.toString('latin1') : `${headerName}` + const lowerName = name.toLowerCase() + + if (lowerName === 'content-encoding' || lowerName === 'content-length') { + continue + } + + filteredHeaders.push(rawHeaders[i], rawHeaders[i + 1]) + } + rawHeaders.splice(0, rawHeaders.length, ...filteredHeaders) + } else if (typeof rawHeaders === 'object') { + for (const name of Object.keys(rawHeaders)) { + const lowerName = name.toLowerCase() + if (lowerName === 'content-encoding' || lowerName === 'content-length') { + delete rawHeaders[name] + } + } + } } - return super.onResponseStart(controller, statusCode, newHeaders, statusMessage) + this.#setupInputBackpressure() + if (this.#decompressors.length === 1) { + this.#setupSingleDecompressor() + } else { + this.#setupMultipleDecompressors() + } + + return super.onResponseStart(this.#controller, statusCode, newHeaders, statusMessage) } /** @@ -45728,10 +46241,13 @@ class DecompressHandler extends DecoratorHandler { */ onResponseData (controller, chunk) { if (this.#decompressors.length > 0) { - this.#decompressors[0].write(chunk) + if (!this.#decompressors[0].write(chunk)) { + this.#inputBackpressured = true + this.#pauseUpstream() + } return } - super.onResponseData(controller, chunk) + return super.onResponseData(this.#controller, chunk) } /** @@ -45741,11 +46257,12 @@ class DecompressHandler extends DecoratorHandler { */ onResponseEnd (controller, trailers) { if (this.#decompressors.length > 0) { + this.#inputEnded = true + this.#trailers = trailers this.#decompressors[0].end() - this.#cleanupDecompressors() return } - super.onResponseEnd(controller, trailers) + return super.onResponseEnd(this.#controller, trailers) } /** @@ -45754,13 +46271,16 @@ class DecompressHandler extends DecoratorHandler { * @returns {void} */ onResponseError (controller, err) { - if (this.#decompressors.length > 0) { - for (const decompressor of this.#decompressors) { - decompressor.destroy(err) - } - this.#cleanupDecompressors() + if (this.#terminated) { + return } - super.onResponseError(controller, err) + + this.#terminated = true + for (const decompressor of this.#decompressors) { + decompressor.destroy() + } + this.#cleanupDecompressors() + super.onResponseError(this.#controller, err) } } @@ -46509,7 +47029,6 @@ class DumpHandler extends DecoratorHandler { #maxSize = 1024 * 1024 #dumped = false #size = 0 - #controller = null aborted = false reason = false @@ -46531,7 +47050,6 @@ class DumpHandler extends DecoratorHandler { onRequestStart (controller, context) { controller.abort = this.#abort.bind(this) - this.#controller = controller return super.onRequestStart(controller, context) } @@ -46555,43 +47073,32 @@ class DumpHandler extends DecoratorHandler { } onResponseError (controller, err) { - if (this.#dumped) { - return - } - - // On network errors before connect, controller will be null - err = this.#controller?.reason ?? err - - super.onResponseError(controller, err) + super.onResponseError(controller, this.aborted === true ? this.reason : err) } onResponseData (controller, chunk) { this.#size = this.#size + chunk.length - if (this.#size >= this.#maxSize) { - this.#dumped = true + if (this.#size > this.#maxSize) { + throw new RequestAbortedError( + `Response size (${this.#size}) larger than maxSize (${this.#maxSize})` + ) + } - if (this.aborted === true) { - super.onResponseError(controller, this.reason) - } else { - super.onResponseEnd(controller, {}) - } + if (this.#size === this.#maxSize) { + this.#dumped = true } return true } onResponseEnd (controller, trailers) { - if (this.#dumped) { - return - } - - if (this.#controller.aborted === true) { + if (this.aborted === true) { super.onResponseError(controller, this.reason) return } - super.onResponseEnd(controller, trailers) + super.onResponseEnd(controller, this.#dumped ? {} : trailers) } } @@ -54018,6 +54525,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} @@ -54058,11 +54608,14 @@ class EventSourceStream extends Transform { eventEndCheck = false /** - * @type {Buffer|null} + * @type {Buffer[]} */ - buffer = null + chunks = [] + chunkIndex = 0 pos = 0 + lineChunkIndex = 0 + linePos = 0 event = { data: undefined, @@ -54102,92 +54655,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) { @@ -54200,10 +54681,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 @@ -54219,19 +54699,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 !== undefined || this.event.retry) { + this.consumeCurrentByte() + if (this.hasPendingEvent()) { this.processEvent(this.event) } this.clearEvent() @@ -54245,22 +54723,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. @@ -54268,7 +54742,7 @@ class EventSourceStream extends Transform { continue } - this.pos++ + this.advanceCursor() } callback() @@ -54293,64 +54767,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 + } } } @@ -54380,13 +54843,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 = { @@ -65234,13 +65836,13 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // The presence of a session property on the socket indicates HTTP2 // HTTP1 if (response.socket?.session == null) { - failWebsocketConnection(handler, 1002, 'Received network error or non-101 status code.', response.error) + failHandshake(handler, response, 1002, 'Received network error or non-101 status code.', response.error) return } // HTTP2 if (response.status !== 200) { - failWebsocketConnection(handler, 1002, 'Received network error or non-200 status code.', response.error) + failHandshake(handler, response, 1002, 'Received network error or non-200 status code.', response.error) return } } @@ -65255,7 +65857,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // header list results in null, failure, or the empty byte // sequence, then fail the WebSocket connection. if (protocols.length !== 0 && !response.headersList.get('Sec-WebSocket-Protocol')) { - failWebsocketConnection(handler, 1002, 'Server did not respond with sent protocols.') + failHandshake(handler, response, 1002, 'Server did not respond with sent protocols.') return } @@ -65271,7 +65873,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // _Fail the WebSocket Connection_. // For H2, no upgrade header is expected. if (response.socket.session == null && response.headersList.get('Upgrade')?.toLowerCase() !== 'websocket') { - failWebsocketConnection(handler, 1002, 'Server did not set Upgrade header to "websocket".') + failHandshake(handler, response, 1002, 'Server did not set Upgrade header to "websocket".') return } @@ -65281,7 +65883,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // MUST _Fail the WebSocket Connection_. // For H2, no connection header is expected. if (response.socket.session == null && response.headersList.get('Connection')?.toLowerCase() !== 'upgrade') { - failWebsocketConnection(handler, 1002, 'Server did not set Connection header to "upgrade".') + failHandshake(handler, response, 1002, 'Server did not set Connection header to "upgrade".') return } @@ -65295,7 +65897,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) const secWSAccept = response.headersList.get('Sec-WebSocket-Accept') const digest = crypto.hash('sha1', keyValue + uid, 'base64') if (secWSAccept !== digest) { - failWebsocketConnection(handler, 1002, 'Incorrect hash received in Sec-WebSocket-Accept header.') + failHandshake(handler, response, 1002, 'Incorrect hash received in Sec-WebSocket-Accept header.') return } @@ -65313,7 +65915,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) extensions = parseExtensions(secExtension) if (!extensions.has('permessage-deflate')) { - failWebsocketConnection(handler, 1002, 'Sec-WebSocket-Extensions header does not match.') + failHandshake(handler, response, 1002, 'Sec-WebSocket-Extensions header does not match.') return } } @@ -65333,8 +65935,8 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // 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)) { - failWebsocketConnection(handler, 1002, 'Protocol was not set in the opening handshake.') + if (requestProtocols === null || !requestProtocols.includes(secProtocol)) { + failHandshake(handler, response, 1002, 'Protocol was not set in the opening handshake.') return } } @@ -65429,6 +66031,16 @@ function closeWebSocketConnection (object, code, reason, validate = false) { } } +function failHandshake (handler, response, code, reason, cause) { + // The H2 upgrade request has already completed and handed off its stream. + // Aborting the request cannot close that stream after handshake validation fails. + if (response.socket?.session != null && !response.socket.destroyed) { + response.socket.destroy() + } + + failWebsocketConnection(handler, code, reason, cause) +} + /** * @param {import('./websocket').Handler} handler * @param {number} code @@ -66147,7 +66759,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 } @@ -66970,9 +67587,9 @@ class WebSocketStream { /** @type {ReadableStreamDefaultController} */ #readableStreamController - // Each WebSocketStream object has an associated writable stream , which is a WritableStream . - /** @type {WritableStream} */ - #writableStream + // Retain the controller so the writable stream can be errored while locked. + /** @type {WritableStreamDefaultController} */ + #writableStreamController // Each WebSocketStream object has an associated boolean handshake aborted , which is initially false. #handshakeAborted = false @@ -67236,6 +67853,9 @@ class WebSocketStream { // 12. Let writable be a new WritableStream . // 13. Set up writable with writeAlgorithm , closeAlgorithm , and abortAlgorithm . const writable = new WritableStream({ + start: (controller) => { + this.#writableStreamController = controller + }, write: (chunk) => this.#write(chunk), close: () => closeWebSocketConnection(this.#handler, null, null), abort: (reason) => this.#closeUsingReason(reason) @@ -67244,9 +67864,6 @@ class WebSocketStream { // Set stream ’s readable stream to readable . this.#readableStream = readable - // Set stream ’s writable stream to writable . - this.#writableStream = writable - // Resolve stream ’s opened promise with WebSocketOpenInfo «[ " extensions " → extensions , " protocol " → protocol , " readable " → readable , " writable " → writable ]». this.#openedPromise.resolve({ extensions, @@ -67332,9 +67949,7 @@ class WebSocketStream { this.#readableStreamController.close() // 6.2. Error stream ’s writable stream with an " InvalidStateError " DOMException indicating that a closed WebSocketStream cannot be written to. - if (!this.#writableStream.locked) { - this.#writableStream.abort(new DOMException('A closed WebSocketStream cannot be written to', 'InvalidStateError')) - } + this.#writableStreamController.error(new DOMException('A closed WebSocketStream cannot be written to', 'InvalidStateError')) // 6.3. Resolve stream ’s closed promise with WebSocketCloseInfo «[ " closeCode " → code , " reason " → reason ]». this.#closedPromise.resolve({ @@ -67351,7 +67966,7 @@ class WebSocketStream { this.#readableStreamController?.error(error) // 7.3. Error stream ’s writable stream with error . - this.#writableStream?.abort(error) + this.#writableStreamController?.error(error) // 7.4. Reject stream ’s closed promise with error . this.#closedPromise.reject(error) diff --git a/dist/run/index.js b/dist/run/index.js index b5a4ce4..6cd357d 100644 --- a/dist/run/index.js +++ b/dist/run/index.js @@ -34858,11 +34858,32 @@ class Request { } } - onUpgrade (statusCode, headers, socket) { + /** + * @param {number} statusCode + * @param {Buffer[]|string[]} 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 (channels.headers.hasSubscribers) { + channels.headers.publish({ request: this, response: { statusCode, headers, statusText } }) + } + + const result = this[kHandler].onUpgrade(statusCode, headers, socket) + + if (!this.aborted) { + this.completed = true + if (channels.trailers.hasSubscribers) { + channels.trailers.publish({ request: this, trailers: [] }) + } + } + + return result } onComplete (trailers) { @@ -37144,14 +37165,16 @@ function defaultFactory (origin, opts) { } class BalancedPool extends PoolBase { - constructor (upstreams = [], { factory = defaultFactory, ...opts } = {}) { + constructor (upstreams = [], { factory = defaultFactory, connect, tls, ...opts } = {}) { if (typeof factory !== 'function') { throw new InvalidArgumentError('factory must be a function.') } super(opts) - this[kOptions] = { ...util.deepClone(opts) } + if (connect && typeof connect !== 'function') connect = { ...connect } + if (tls && typeof tls !== 'function') tls = { ...tls } + this[kOptions] = { ...util.deepClone(opts), connect, tls } this[kOptions].interceptors = opts.interceptors ? { ...opts.interceptors } : undefined @@ -37391,8 +37414,14 @@ function lazyllhttp () { let mod - // We disable wasm SIMD on ppc64 as it seems to be broken on Power 9 architectures. - let useWasmSIMD = process.arch !== 'ppc64' + // We disable wasm SIMD on older versions of Node.js on ppc64 that are broken on Power >=9 architectures. + let useWasmSIMD = true + if (process.arch === 'ppc64') { + const [major, minor] = process.versions.node.split('.').map(n => parseInt(n, 10)) + if (major < 24 || (major === 24 && minor < 12)) { + useWasmSIMD = false + } + } // The Env Variable UNDICI_NO_WASM_SIMD allows explicitly overriding the default behavior if (process.env.UNDICI_NO_WASM_SIMD === '1') { useWasmSIMD = false @@ -37843,7 +37872,7 @@ class Parser { * @param {Buffer} head */ onUpgrade (head) { - const { upgrade, client, socket, headers, statusCode } = this + const { upgrade, client, socket, headers, statusCode, statusText } = this assert(upgrade) assert(client[kSocket] === socket) @@ -37878,8 +37907,9 @@ class Parser { client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade')) try { - request.onUpgrade(statusCode, headers, socket) + request.onUpgrade(statusCode, headers, socket, statusText) } catch (err) { + util.errorRequest(client, request, err) util.destroy(socket, err) } @@ -38334,7 +38364,7 @@ function onSocketClose () { function clearIdleSocketValidation (socket) { if (socket[kIdleSocketValidationTimeout]) { - clearTimeout(socket[kIdleSocketValidationTimeout]) + clearImmediate(socket[kIdleSocketValidationTimeout]) socket[kIdleSocketValidationTimeout] = null } @@ -38343,15 +38373,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?.() + }) } /** @@ -39069,7 +39107,9 @@ const { RequestAbortedError, SocketError, InformationalError, - InvalidArgumentError + InvalidArgumentError, + HeadersTimeoutError, + BodyTimeoutError } = __nccwpck_require__(68707) const { kUrl, @@ -39094,6 +39134,7 @@ const { kHTTPContext, kClosed, kBodyTimeout, + kHeadersTimeout, kEnableConnectProtocol, kRemoteSettings, kHTTP2Stream, @@ -39280,7 +39321,11 @@ function resumeH2 (client) { const socket = client[kSocket] if (socket?.destroyed === false) { - if (client[kSize] === 0 || client[kMaxConcurrentStreams] === 0) { + // Only let the process exit when there is genuinely nothing outstanding. + // Unreffing because the peer advertised MAX_CONCURRENT_STREAMS = 0 left + // queued requests with nothing holding the event loop open, so the process + // could exit with status 0 while an awaited request never settled. + if (client[kSize] === 0) { socket.unref() client[kHTTP2Session].unref() } else { @@ -39375,6 +39420,36 @@ function onHttp2SessionEnd () { * @this {import('http2').ClientHttp2Session} * @param {number} errorCode */ +// Backport of #5410 and #5569. HTTP/2 multiplexes, so requests complete out of +// order; advancing kRunningIdx blindly retired whichever request happened to +// sit at the head instead of the one that actually finished, which both lost +// requests and left phantom running slots behind. +function completeRequest (client, request, resetPendingIdx = false) { + const queue = client[kQueue] + const runningIdx = client[kRunningIdx] + + // In-order completion: clear the request and advance without splicing. + // The client's resume loop compacts cleared slots once the index grows. + if (runningIdx < client[kPendingIdx] && queue[runningIdx] === request) { + queue[runningIdx] = null + client[kRunningIdx] = runningIdx + 1 + return + } + + const index = queue.indexOf(request, runningIdx) + + if (index === -1 || index >= client[kPendingIdx]) { + return + } + + queue.splice(index, 1) + client[kPendingIdx]-- + + if (resetPendingIdx && client[kPendingIdx] < client[kRunningIdx]) { + client[kPendingIdx] = client[kRunningIdx] + } +} + function onHttp2SessionGoAway (errorCode) { // TODO(mcollina): Verify if GOAWAY implements the spec correctly: // https://datatracker.ietf.org/doc/html/rfc7540#section-6.8 @@ -39396,7 +39471,9 @@ function onHttp2SessionGoAway (errorCode) { if (client[kRunningIdx] < client[kQueue].length) { const request = client[kQueue][client[kRunningIdx]] client[kQueue][client[kRunningIdx]++] = null - util.errorRequest(client, request, err) + if (request != null) { + util.errorRequest(client, request, err) + } client[kPendingIdx] = client[kRunningIdx] } @@ -39429,7 +39506,9 @@ function onHttp2SessionClose () { const requests = client[kQueue].splice(client[kRunningIdx]) for (let i = 0; i < requests.length; i++) { const request = requests[i] - util.errorRequest(client, request, err) + if (request != null) { + util.errorRequest(client, request, err) + } } } } @@ -39477,7 +39556,10 @@ function shouldSendContentLength (method) { } function writeH2 (client, request) { - const requestTimeout = request.bodyTimeout ?? client[kBodyTimeout] + // Time to the response headers, then time between body chunks. Using + // bodyTimeout for both made headersTimeout a no-op over HTTP/2. + const headersTimeout = request.headersTimeout ?? client[kHeadersTimeout] + const bodyTimeout = request.bodyTimeout ?? client[kBodyTimeout] const session = client[kHTTP2Session] const { method, path, host, upgrade, expectContinue, signal, protocol, headers: reqHeaders } = request let { body } = request @@ -39544,6 +39626,7 @@ function writeH2 (client, request) { // We move the running index to the next request client[kOnError](err) + completeRequest(client, request) client[kResume]() } @@ -39591,14 +39674,18 @@ function writeH2 (client, request) { stream = session.request(headers, { endStream: false, signal }) stream[kHTTP2Stream] = true + ++session[kOpenStreams] stream.once('response', (headers, _flags) => { const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers - request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) - - ++session[kOpenStreams] - client[kQueue][client[kRunningIdx]++] = null + try { + request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) + } catch (err) { + abort(err) + return + } + completeRequest(client, request) }) stream.on('error', () => { @@ -39615,7 +39702,7 @@ function writeH2 (client, request) { if (session[kOpenStreams] === 0) session.unref() }) - stream.setTimeout(requestTimeout) + stream.setTimeout(headersTimeout) return true } @@ -39626,18 +39713,24 @@ function writeH2 (client, request) { // We disabled endStream to allow the user to write to the stream stream = session.request(headers, { endStream: false, signal }) stream[kHTTP2Stream] = true + ++session[kOpenStreams] stream.on('response', headers => { const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers - request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) - ++session[kOpenStreams] - client[kQueue][client[kRunningIdx]++] = null + try { + request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream) + } catch (err) { + abort(err) + return + } + completeRequest(client, request) }) + stream.on('error', abort) stream.once('close', () => { session[kOpenStreams] -= 1 if (session[kOpenStreams] === 0) session.unref() }) - stream.setTimeout(requestTimeout) + stream.setTimeout(headersTimeout) return true } @@ -39738,7 +39831,7 @@ function writeH2 (client, request) { // Increment counter as we have new streams open ++session[kOpenStreams] - stream.setTimeout(requestTimeout) + stream.setTimeout(headersTimeout) // Track whether we received a response (headers) let responseReceived = false @@ -39747,6 +39840,7 @@ function writeH2 (client, request) { const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers request.onResponseStarted() responseReceived = true + stream.setTimeout(bodyTimeout) // Due to the stream nature, it is possible we face a race condition // where the stream has been assigned, but the request has been aborted @@ -39781,14 +39875,13 @@ function writeH2 (client, request) { request.onComplete({}) } - client[kQueue][client[kRunningIdx]++] = null + completeRequest(client, request) client[kResume]() } else { // Stream ended without receiving a response - this is an error // (e.g., server destroyed the stream before sending headers) abort(new InformationalError('HTTP/2: stream half-closed (remote)')) - client[kQueue][client[kRunningIdx]++] = null - client[kPendingIdx] = client[kRunningIdx] + completeRequest(client, request, true) client[kResume]() } }) @@ -39799,6 +39892,14 @@ function writeH2 (client, request) { if (session[kOpenStreams] === 0) { session.unref() } + + // A stream can close without ever emitting 'end' or 'error': a peer's + // RST_STREAM(CANCEL) received before the response is reported by Node as a + // bare 'close', and destroying the stream unenrolls its timeout, so no + // 'timeout' follows either. Nothing else would ever settle this request. + if (!request.aborted && !request.completed) { + abort(new InformationalError('HTTP/2: stream closed before the response was complete')) + } }) stream.once('error', function (err) { @@ -39816,7 +39917,9 @@ function writeH2 (client, request) { }) stream.on('timeout', () => { - const err = new InformationalError(`HTTP/2: "stream timeout after ${requestTimeout}"`) + const err = responseReceived + ? new BodyTimeoutError(`HTTP/2: "body timeout after ${bodyTimeout}"`) + : new HeadersTimeoutError(`HTTP/2: "headers timeout after ${headersTimeout}"`) stream.removeAllListeners('data') session[kOpenStreams] -= 1 @@ -40438,7 +40541,9 @@ class Client extends DispatcherBase { const requests = this[kQueue].splice(this[kPendingIdx]) for (let i = 0; i < requests.length; i++) { const request = requests[i] - util.errorRequest(this, request, err) + if (request != null) { + util.errorRequest(this, request, err) + } } const callback = () => { @@ -40477,7 +40582,9 @@ function onError (client, err) { for (let i = 0; i < requests.length; i++) { const request = requests[i] - util.errorRequest(client, request, err) + if (request != null) { + util.errorRequest(client, request, err) + } } assert(client[kSize] === 0) } @@ -42786,6 +42893,13 @@ class CacheHandler { } const cacheControlHeader = resHeaders['cache-control'] + const cacheControlDirectives = cacheControlHeader ? parseCacheControlHeader(cacheControlHeader) : {} + + if (revalidationResponseDisallowsCachedReuse(this.#cacheType, resHeaders, cacheControlDirectives)) { + deleteCachedValue(this.#store, this.#cacheKey) + return downstreamOnHeaders() + } + const heuristicallyCacheable = resHeaders['last-modified'] && arrayIncludes(HEURISTICALLY_CACHEABLE_STATUS_CODES, statusCode) if ( !cacheControlHeader && @@ -42802,8 +42916,7 @@ class CacheHandler { return downstreamOnHeaders() } - const cacheControlDirectives = cacheControlHeader ? parseCacheControlHeader(cacheControlHeader) : {} - if (!canCacheResponse(this.#cacheType, statusCode, resHeaders, cacheControlDirectives, this.#cacheKey.headers)) { + if (!canCacheResponse(this.#cacheType, this.#cacheKey.method, statusCode, resHeaders, cacheControlDirectives, this.#cacheKey.headers)) { if (statusCode === 304 && (cacheControlHeader || revalidationResponseDisallowsCachedReuse(this.#cacheType, resHeaders, cacheControlDirectives))) { deleteCachedValue(this.#store, this.#cacheKey) } @@ -43044,7 +43157,10 @@ function deleteCachedValueIfNotModified (statusCode, store, cacheKey) { */ function revalidationResponseDisallowsCachedReuse (cacheType, resHeaders, cacheControlDirectives) { return cacheControlDirectives['no-store'] === true || - (cacheType === 'shared' && cacheControlDirectives.private === true) || + (cacheType === 'shared' && ( + cacheControlDirectives.private === true || + Object.hasOwn(resHeaders, 'set-cookie') + )) || (resHeaders.vary ? isInvalidOrWildcardVaryHeader(resHeaders.vary) : false) } @@ -43052,12 +43168,16 @@ function revalidationResponseDisallowsCachedReuse (cacheType, resHeaders, cacheC * @see https://www.rfc-editor.org/rfc/rfc9111.html#name-storing-responses-to-authen * * @param {import('../../types/cache-interceptor.d.ts').default.CacheOptions['type']} cacheType + * @param {string} method * @param {number} statusCode * @param {import('../../types/header.d.ts').IncomingHttpHeaders} resHeaders * @param {import('../../types/cache-interceptor.d.ts').default.CacheControlDirectives} cacheControlDirectives * @param {import('../../types/header.d.ts').IncomingHttpHeaders} [reqHeaders] */ -function canCacheResponse (cacheType, statusCode, resHeaders, cacheControlDirectives, reqHeaders) { +function canCacheResponse (cacheType, method, statusCode, resHeaders, cacheControlDirectives, reqHeaders) { + if (!arrayIncludes(util.safeHTTPMethods, method)) { + return false + } // Status code must be final and understood. if (statusCode < 200 || arrayIncludes(NOT_UNDERSTOOD_STATUS_CODES, statusCode)) { return false @@ -43078,7 +43198,10 @@ function canCacheResponse (cacheType, statusCode, resHeaders, cacheControlDirect return false } - if (cacheType === 'shared' && cacheControlDirectives.private === true) { + if (cacheType === 'shared' && ( + cacheControlDirectives.private === true || + Object.hasOwn(resHeaders, 'set-cookie') + )) { return false } @@ -44303,6 +44426,55 @@ function validatePartialResponseContentLength (headers, range, statusCode, retry } } +// A stable controller handed to the downstream handler for the lifetime of the +// request. Each transparent retry/resume is a separate dispatch with its own +// connection controller. The proxy always forwards to the active connection +// while preserving a downstream pause across controller replacement. +class RetryController { + #paused = false + #target = null + + set target (target) { + this.#target = target + if (this.#paused) { + target?.pause() + } + } + + get target () { return this.#target } + + pause () { + this.#paused = true + this.#target?.pause() + } + + resume () { + this.#paused = false + this.#target?.resume() + } + + abort (reason) { + this.#target?.abort(reason) + } + + get paused () { return this.#paused || (this.#target?.paused ?? false) } + get aborted () { return this.#target?.aborted ?? false } + get reason () { return this.#target?.reason ?? null } + get rawHeaders () { return this.#target?.rawHeaders ?? null } + set rawHeaders (value) { + if (this.#target) { + this.#target.rawHeaders = value + } + } + + get rawTrailers () { return this.#target?.rawTrailers ?? null } + set rawTrailers (value) { + if (this.#target) { + this.#target.rawTrailers = value + } + } +} + class RetryHandler { constructor (opts, { dispatch, handler }) { const { retryOptions, ...dispatchOpts } = opts @@ -44357,14 +44529,23 @@ class RetryHandler { this.start = 0 this.end = null this.etag = null + this.controllerProxy = new RetryController() } onResponseStartWithRetry (controller, statusCode, headers, statusMessage, err) { if (this.retryOpts.throwOnError) { // Preserve old behavior for status codes that are not eligible for retry if (this.retryOpts.statusCodes.includes(statusCode) === false) { - this.headersSent = true - this.handler.onResponseStart?.(controller, statusCode, headers, statusMessage) + if (this.headersSent) { + // The downstream handler already received the response from an + // earlier attempt. Forwarding this response would replace the + // downstream body and leave the original body pending forever. + this.handler.onResponseError?.(this.controllerProxy, err) + } else { + this.headersSent = true + this.checkpointResponseEnd(headers) + this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage) + } } else { this.error = err } @@ -44374,14 +44555,23 @@ class RetryHandler { if (isDisturbed(this.opts.body)) { this.headersSent = true - this.handler.onResponseStart?.(controller, statusCode, headers, statusMessage) + this.checkpointResponseEnd(headers) + this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage) return } function shouldRetry (passedErr) { if (passedErr) { - this.headersSent = true - this.handler.onResponseStart?.(controller, statusCode, headers, statusMessage) + if (this.headersSent) { + // The downstream handler already received the response from an + // earlier attempt. Forwarding this response would replace the + // downstream body and leave the original body pending forever. + this.handler.onResponseError?.(this.controllerProxy, passedErr) + } else { + this.headersSent = true + this.checkpointResponseEnd(headers) + this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage) + } controller.resume() return } @@ -44401,14 +44591,29 @@ class RetryHandler { ) } - onRequestStart (controller, context) { - if (!this.headersSent) { - this.handler.onRequestStart?.(controller, context) + checkpointResponseEnd (headers) { + 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 } } - onRequestUpgrade (controller, statusCode, headers, socket) { - this.handler.onRequestUpgrade?.(controller, statusCode, headers, socket) + onRequestStart (controller, context) { + this.controllerProxy.target = controller + if (!this.headersSent) { + this.handler.onRequestStart?.(this.controllerProxy, context) + } + } + + onRequestUpgrade (_controller, statusCode, headers, socket) { + this.handler.onRequestUpgrade?.(this.controllerProxy, statusCode, headers, socket) } static [kRetryHandlerDefaultRetry] (err, { state, opts }, cb) { @@ -44521,8 +44726,12 @@ class RetryHandler { const { start, size, end = size ? size - 1 : null } = 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)) { + throw new RequestRetryError('Content-Range mismatch', statusCode, { + headers, + data: { count: this.retryCount } + }) + } return } @@ -44535,7 +44744,7 @@ class RetryHandler { if (range == null) { this.headersSent = true this.handler.onResponseStart?.( - controller, + this.controllerProxy, statusCode, headers, statusMessage @@ -44584,7 +44793,7 @@ class RetryHandler { this.headersSent = true this.handler.onResponseStart?.( - controller, + this.controllerProxy, statusCode, headers, statusMessage @@ -44597,30 +44806,30 @@ class RetryHandler { } } - onResponseData (controller, chunk) { + onResponseData (_controller, chunk) { if (this.error) { return } this.start += chunk.length - this.handler.onResponseData?.(controller, chunk) + this.handler.onResponseData?.(this.controllerProxy, chunk) } - onResponseEnd (controller, trailers) { + onResponseEnd (_controller, trailers) { if (this.error && this.retryOpts.throwOnError) { throw this.error } if (!this.error) { this.retryCount = 0 - return this.handler.onResponseEnd?.(controller, trailers) + return this.handler.onResponseEnd?.(this.controllerProxy, trailers) } - this.retry(controller) + this.retry() } - retry (controller) { + retry () { if (this.start !== 0) { const headers = { range: `bytes=${this.start}-${this.end ?? ''}` } @@ -44642,23 +44851,23 @@ class RetryHandler { this.retryCountCheckpoint = this.retryCount this.dispatch(this.opts, this) } catch (err) { - this.handler.onResponseError?.(controller, err) + this.handler.onResponseError?.(this.controllerProxy, err) } } onResponseError (controller, err) { - if (controller?.aborted || isDisturbed(this.opts.body)) { - this.handler.onResponseError?.(controller, err) + if (controller?.aborted || isDisturbed(this.opts.body) || (this.headersSent && !this.resume)) { + this.handler.onResponseError?.(this.controllerProxy, err) return } function shouldRetry (returnedErr) { if (!returnedErr) { - this.retry(controller) + this.retry() return } - this.handler?.onResponseError?.(controller, returnedErr) + this.handler?.onResponseError?.(this.controllerProxy, returnedErr) } // We reconcile in case of a mix between network errors @@ -45038,7 +45247,10 @@ function staleResponseRequiresRevalidation (result, cacheType) { * @returns {boolean} */ function revalidationResponseDisallowsCachedReuse (cacheType, headers) { - if (headers.vary && isInvalidOrWildcardVaryHeader(headers.vary)) { + if ( + (headers.vary && isInvalidOrWildcardVaryHeader(headers.vary)) || + (cacheType === 'shared' && Object.hasOwn(headers, 'set-cookie')) + ) { return true } @@ -45297,6 +45509,17 @@ function handleResult ( return handleUncachedResponse(dispatch, globalOpts, cacheKey, handler, opts, reqCacheControl) } + // Shared stores may outlive the Undici version that wrote them. Do not + // re-serve a Set-Cookie header from an existing shared-cache entry. + if (globalOpts.type === 'shared' && Object.hasOwn(result.headers, 'set-cookie')) { + if (util.isStream(result.body)) { + result.body.on('error', nop).destroy() + } + + deleteCachedValue(globalOpts.store, cacheKey) + return handleUncachedResponse(dispatch, globalOpts, cacheKey, handler, opts, reqCacheControl) + } + const now = Date.now() if (now > result.deleteAt) { // Response is expired, cache store shouldn't have given this to us @@ -45495,6 +45718,11 @@ module.exports = (opts = {}) => { * @type {import('../../types/cache-interceptor.d.ts').default.CacheKey} */ const cacheKey = makeCacheKey(opts) + + if (!arrayIncludes(util.safeHTTPMethods, opts.method)) { + return dispatch(opts, new CacheHandler(globalOpts, cacheKey, handler)) + } + const result = store.get(cacheKey) if (result && typeof result.then === 'function') { @@ -45532,7 +45760,8 @@ module.exports = (opts = {}) => { const { createInflate, createGunzip, createBrotliDecompress, createZstdDecompress } = __nccwpck_require__(38522) -const { pipeline } = __nccwpck_require__(57075) +const { pipeline, Transform: TransformStream } = __nccwpck_require__(57075) +const { InvalidArgumentError, ResponseExceededMaxSizeError } = __nccwpck_require__(68707) const DecoratorHandler = __nccwpck_require__(58155) const { runtimeFeatures } = __nccwpck_require__(313) @@ -45540,6 +45769,60 @@ const { runtimeFeatures } = __nccwpck_require__(313) /** @typedef {import('node:stream').Transform} Controller */ /** @typedef {Transform&import('node:zlib').Zlib} DecompressorStream */ +class DecompressController { + #onPause + #onResume + #onAbort + #paused = false + + constructor (onPause, onResume, onAbort) { + this.#onPause = onPause + this.#onResume = onResume + this.#onAbort = onAbort + this.target = null + } + + pause () { + if (this.#paused) { + return + } + + this.#paused = true + this.#onPause() + } + + resume () { + if (!this.#paused) { + return + } + + this.#paused = false + this.#onResume() + } + + abort (reason) { + this.target?.abort(reason) + this.#onAbort(reason) + } + + get paused () { return this.#paused } + get aborted () { return this.target?.aborted ?? false } + get reason () { return this.target?.reason ?? null } + get rawHeaders () { return this.target?.rawHeaders ?? null } + set rawHeaders (value) { + if (this.target) { + this.target.rawHeaders = value + } + } + + get rawTrailers () { return this.target?.rawTrailers ?? null } + set rawTrailers (value) { + if (this.target) { + this.target.rawTrailers = value + } + } +} + /** @type {Record DecompressorStream>} */ const supportedEncodings = { gzip: createGunzip, @@ -45552,6 +45835,31 @@ const supportedEncodings = { } const defaultSkipStatusCodes = /** @type {const} */ ([204, 304]) +const defaultMaxSize = 0 + +/** + * Limits the output of one stage in a decompression chain. + * @param {number} maxSize - Maximum output size in bytes + * @returns {Transform} + */ +function createMaxSizeLimiter (maxSize) { + let size = 0 + + return new TransformStream({ + transform (chunk, _encoding, callback) { + const decompressedSize = size + chunk.length + if (decompressedSize > maxSize) { + callback(new ResponseExceededMaxSizeError( + `Decompressed response size (${decompressedSize}) exceeded maxSize (${maxSize})` + )) + return + } + + size = decompressedSize + callback(null, chunk) + } + }) +} let warningEmitted = /** @type {boolean} */ (false) @@ -45559,20 +45867,150 @@ let warningEmitted = /** @type {boolean} */ (false) * @typedef {Object} DecompressHandlerOptions * @property {number[]|Readonly} [skipStatusCodes=[204, 304]] - List of status codes to skip decompression for * @property {boolean} [skipErrorResponses] - Whether to skip decompression for error responses (status codes >= 400) + * @property {number} [maxSize=0] - Maximum decompressed response size in bytes. 0 disables the limit */ class DecompressHandler extends DecoratorHandler { /** @type {Transform[]} */ #decompressors = [] + /** @type {Record | undefined} */ + #trailers /** @type {Readonly} */ #skipStatusCodes /** @type {boolean} */ #skipErrorResponses + /** @type {number} */ + #maxSize + /** @type {number} */ + #decompressedSize = 0 + /** @type {boolean} */ + #terminated = false + /** @type {boolean} */ + #inputEnded = false + /** @type {boolean} */ + #inputBackpressured = false + /** @type {boolean} */ + #upstreamPaused = false + /** @type {boolean} */ + #draining = false + /** @type {boolean} */ + #drainRequested = false + /** @type {boolean} */ + #completionPending = false + /** @type {DecompressorStream | undefined} */ + #finalDecompressor + /** @type {DecompressController} */ + #controller + + constructor (handler, { skipStatusCodes = defaultSkipStatusCodes, skipErrorResponses = true, maxSize = defaultMaxSize } = {}) { + if (!Number.isSafeInteger(maxSize) || maxSize < 0) { + throw new InvalidArgumentError('maxSize must be a non-negative integer') + } - constructor (handler, { skipStatusCodes = defaultSkipStatusCodes, skipErrorResponses = true } = {}) { super(handler) this.#skipStatusCodes = skipStatusCodes this.#skipErrorResponses = skipErrorResponses + this.#maxSize = maxSize + this.#controller = new DecompressController( + () => this.#onDownstreamPause(), + () => this.#onDownstreamResume(), + reason => { + if (this.#inputEnded && !this.#terminated) { + this.onResponseError(this.#controller, reason) + } + } + ) + } + + #onDownstreamPause () { + this.#pauseUpstream() + } + + #onDownstreamResume () { + const drainWasDeferred = this.#draining + this.#drainOutput() + if (!drainWasDeferred) { + this.#resumeUpstreamIfNeeded() + this.#finishIfReady() + } + } + + #pauseUpstream () { + if (!this.#upstreamPaused && !this.#terminated) { + this.#upstreamPaused = true + this.#controller.target?.pause() + } + } + + #resumeUpstreamIfNeeded () { + if (this.#upstreamPaused && !this.#controller.paused && !this.#inputBackpressured) { + this.#upstreamPaused = false + if (!this.#inputEnded) { + this.#controller.target?.resume() + } + } + } + + #drainOutput () { + if (this.#terminated || this.#controller.paused || !this.#finalDecompressor) { + return + } + + if (this.#draining) { + this.#drainRequested = true + return + } + + this.#draining = true + try { + do { + this.#drainRequested = false + let chunk + while (!this.#terminated && !this.#controller.paused && (chunk = this.#finalDecompressor.read()) !== null) { + if (this.#maxSize > 0) { + const decompressedSize = this.#decompressedSize + chunk.length + if (decompressedSize > this.#maxSize) { + this.#fail(new ResponseExceededMaxSizeError( + `Decompressed response size (${decompressedSize}) exceeded maxSize (${this.#maxSize})` + )) + return + } + + this.#decompressedSize = decompressedSize + } + + const result = super.onResponseData(this.#controller, chunk) + if (result === false && !this.#controller.paused) { + this.#controller.pause() + } + } + } while (this.#drainRequested && !this.#terminated && !this.#controller.paused) + } finally { + this.#draining = false + } + + this.#resumeUpstreamIfNeeded() + this.#finishIfReady() + } + + #finishIfReady () { + if (this.#terminated || !this.#completionPending || this.#controller.paused || this.#draining) { + return + } + + this.#terminated = true + this.#cleanupDecompressors() + super.onResponseEnd(this.#controller, this.#trailers) + } + + #onDecompressionEnd () { + if (this.#terminated) { + return + } + + this.#completionPending = true + this.#drainOutput() + this.#finishIfReady() } /** @@ -45592,7 +46030,7 @@ class DecompressHandler extends DecoratorHandler { * Creates a chain of decompressors for multiple content encodings * * @param {string} encodings - Comma-separated list of content encodings - * @returns {Array} - Array of decompressor streams + * @returns {Array} - Array of decompressor and limiting streams * @throws {Error} - If the number of content-encodings exceeds the maximum allowed */ #createDecompressionChain (encodings) { @@ -45620,60 +46058,97 @@ class DecompressHandler extends DecoratorHandler { decompressors.push(supportedEncodings[encoding]()) } - return decompressors + if (decompressors.length < 2) { + return decompressors + } + + /** @type {Transform[]} */ + const streams = [] + for (let i = 0; i < decompressors.length; i++) { + streams.push(decompressors[i]) + if (i < decompressors.length - 1 && this.#maxSize > 0) { + streams.push(createMaxSizeLimiter(this.#maxSize)) + } + } + + return streams } /** - * Sets up event handlers for a decompressor stream using readable events - * @param {DecompressorStream} decompressor - The decompressor stream - * @param {Controller} controller - The controller to coordinate with + * Stops decompression and reports an error. + * @param {Error} error - The decompression error * @returns {void} */ - #setupDecompressorEvents (decompressor, controller) { - decompressor.on('readable', () => { - let chunk - while ((chunk = decompressor.read()) !== null) { - const result = super.onResponseData(controller, chunk) - if (result === false) { - break - } - } - }) + #fail (error) { + if (this.#terminated) { + return + } - decompressor.on('error', (error) => { - super.onResponseError(controller, error) - }) + if (this.#inputEnded) { + // The request is already marked complete once the compressed input ends, + // so controller.abort() can no longer propagate decoder flush errors. + this.onResponseError(this.#controller, error) + } else { + this.#controller.abort(error) + } + } + + /** + * Sets up event handlers for the final decompressor stream. + * @param {DecompressorStream} decompressor - The decompressor stream + * @returns {void} + */ + #setupDecompressorEvents (decompressor) { + this.#finalDecompressor = decompressor + decompressor.on('readable', () => this.#drainOutput()) + decompressor.on('error', (error) => this.#fail(error)) } /** * Sets up event handling for a single decompressor - * @param {Controller} controller - The controller to handle events * @returns {void} */ - #setupSingleDecompressor (controller) { + #setupSingleDecompressor () { const decompressor = this.#decompressors[0] - this.#setupDecompressorEvents(decompressor, controller) + this.#setupDecompressorEvents(decompressor) - decompressor.on('end', () => { - super.onResponseEnd(controller, {}) - }) + decompressor.on('end', () => this.#onDecompressionEnd()) } /** * Sets up event handling for multiple chained decompressors using pipeline - * @param {Controller} controller - The controller to handle events * @returns {void} */ - #setupMultipleDecompressors (controller) { + #setupMultipleDecompressors () { const lastDecompressor = this.#decompressors[this.#decompressors.length - 1] - this.#setupDecompressorEvents(lastDecompressor, controller) + this.#setupDecompressorEvents(lastDecompressor) pipeline(this.#decompressors, (err) => { - if (err) { - super.onResponseError(controller, err) + if (this.#terminated) { return } - super.onResponseEnd(controller, {}) + + if (err) { + this.#fail(err) + return + } + + this.#onDecompressionEnd() + }) + } + + #setupInputBackpressure () { + const decompressor = this.#decompressors[0] + decompressor.on('drain', () => { + if (this.#terminated) { + return + } + + this.#inputBackpressured = false + if (!this.#controller.paused) { + this.#drainOutput() + this.#resumeUpstreamIfNeeded() + } }) } @@ -45683,6 +46158,16 @@ class DecompressHandler extends DecoratorHandler { */ #cleanupDecompressors () { this.#decompressors.length = 0 + this.#finalDecompressor = undefined + } + + onRequestStart (controller, context) { + this.#controller.target = controller + return super.onRequestStart(this.#controller, context) + } + + onRequestUpgrade (controller, statusCode, headers, socket) { + return super.onRequestUpgrade(this.#controller, statusCode, headers, socket) } /** @@ -45697,14 +46182,14 @@ class DecompressHandler extends DecoratorHandler { // If content encoding is not supported or status code is in skip list if (this.#shouldSkipDecompression(contentEncoding, statusCode)) { - return super.onResponseStart(controller, statusCode, headers, statusMessage) + return super.onResponseStart(this.#controller, statusCode, headers, statusMessage) } const decompressors = this.#createDecompressionChain(contentEncoding.toLowerCase()) if (decompressors.length === 0) { this.#cleanupDecompressors() - return super.onResponseStart(controller, statusCode, headers, statusMessage) + return super.onResponseStart(this.#controller, statusCode, headers, statusMessage) } this.#decompressors = decompressors @@ -45712,13 +46197,41 @@ class DecompressHandler extends DecoratorHandler { // Remove compression headers since we're decompressing const { 'content-encoding': _, 'content-length': __, ...newHeaders } = headers - if (this.#decompressors.length === 1) { - this.#setupSingleDecompressor(controller) - } else { - this.#setupMultipleDecompressors(controller) + if (this.#controller.rawHeaders) { + const rawHeaders = this.#controller.rawHeaders + + if (Array.isArray(rawHeaders)) { + const filteredHeaders = [] + for (let i = 0; i < rawHeaders.length; i += 2) { + const headerName = rawHeaders[i] + const name = Buffer.isBuffer(headerName) ? headerName.toString('latin1') : `${headerName}` + const lowerName = name.toLowerCase() + + if (lowerName === 'content-encoding' || lowerName === 'content-length') { + continue + } + + filteredHeaders.push(rawHeaders[i], rawHeaders[i + 1]) + } + rawHeaders.splice(0, rawHeaders.length, ...filteredHeaders) + } else if (typeof rawHeaders === 'object') { + for (const name of Object.keys(rawHeaders)) { + const lowerName = name.toLowerCase() + if (lowerName === 'content-encoding' || lowerName === 'content-length') { + delete rawHeaders[name] + } + } + } } - return super.onResponseStart(controller, statusCode, newHeaders, statusMessage) + this.#setupInputBackpressure() + if (this.#decompressors.length === 1) { + this.#setupSingleDecompressor() + } else { + this.#setupMultipleDecompressors() + } + + return super.onResponseStart(this.#controller, statusCode, newHeaders, statusMessage) } /** @@ -45728,10 +46241,13 @@ class DecompressHandler extends DecoratorHandler { */ onResponseData (controller, chunk) { if (this.#decompressors.length > 0) { - this.#decompressors[0].write(chunk) + if (!this.#decompressors[0].write(chunk)) { + this.#inputBackpressured = true + this.#pauseUpstream() + } return } - super.onResponseData(controller, chunk) + return super.onResponseData(this.#controller, chunk) } /** @@ -45741,11 +46257,12 @@ class DecompressHandler extends DecoratorHandler { */ onResponseEnd (controller, trailers) { if (this.#decompressors.length > 0) { + this.#inputEnded = true + this.#trailers = trailers this.#decompressors[0].end() - this.#cleanupDecompressors() return } - super.onResponseEnd(controller, trailers) + return super.onResponseEnd(this.#controller, trailers) } /** @@ -45754,13 +46271,16 @@ class DecompressHandler extends DecoratorHandler { * @returns {void} */ onResponseError (controller, err) { - if (this.#decompressors.length > 0) { - for (const decompressor of this.#decompressors) { - decompressor.destroy(err) - } - this.#cleanupDecompressors() + if (this.#terminated) { + return } - super.onResponseError(controller, err) + + this.#terminated = true + for (const decompressor of this.#decompressors) { + decompressor.destroy() + } + this.#cleanupDecompressors() + super.onResponseError(this.#controller, err) } } @@ -46509,7 +47029,6 @@ class DumpHandler extends DecoratorHandler { #maxSize = 1024 * 1024 #dumped = false #size = 0 - #controller = null aborted = false reason = false @@ -46531,7 +47050,6 @@ class DumpHandler extends DecoratorHandler { onRequestStart (controller, context) { controller.abort = this.#abort.bind(this) - this.#controller = controller return super.onRequestStart(controller, context) } @@ -46555,43 +47073,32 @@ class DumpHandler extends DecoratorHandler { } onResponseError (controller, err) { - if (this.#dumped) { - return - } - - // On network errors before connect, controller will be null - err = this.#controller?.reason ?? err - - super.onResponseError(controller, err) + super.onResponseError(controller, this.aborted === true ? this.reason : err) } onResponseData (controller, chunk) { this.#size = this.#size + chunk.length - if (this.#size >= this.#maxSize) { - this.#dumped = true + if (this.#size > this.#maxSize) { + throw new RequestAbortedError( + `Response size (${this.#size}) larger than maxSize (${this.#maxSize})` + ) + } - if (this.aborted === true) { - super.onResponseError(controller, this.reason) - } else { - super.onResponseEnd(controller, {}) - } + if (this.#size === this.#maxSize) { + this.#dumped = true } return true } onResponseEnd (controller, trailers) { - if (this.#dumped) { - return - } - - if (this.#controller.aborted === true) { + if (this.aborted === true) { super.onResponseError(controller, this.reason) return } - super.onResponseEnd(controller, trailers) + super.onResponseEnd(controller, this.#dumped ? {} : trailers) } } @@ -54018,6 +54525,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} @@ -54058,11 +54608,14 @@ class EventSourceStream extends Transform { eventEndCheck = false /** - * @type {Buffer|null} + * @type {Buffer[]} */ - buffer = null + chunks = [] + chunkIndex = 0 pos = 0 + lineChunkIndex = 0 + linePos = 0 event = { data: undefined, @@ -54102,92 +54655,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) { @@ -54200,10 +54681,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 @@ -54219,19 +54699,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 !== undefined || this.event.retry) { + this.consumeCurrentByte() + if (this.hasPendingEvent()) { this.processEvent(this.event) } this.clearEvent() @@ -54245,22 +54723,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. @@ -54268,7 +54742,7 @@ class EventSourceStream extends Transform { continue } - this.pos++ + this.advanceCursor() } callback() @@ -54293,64 +54767,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 + } } } @@ -54380,13 +54843,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 = { @@ -65234,13 +65836,13 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // The presence of a session property on the socket indicates HTTP2 // HTTP1 if (response.socket?.session == null) { - failWebsocketConnection(handler, 1002, 'Received network error or non-101 status code.', response.error) + failHandshake(handler, response, 1002, 'Received network error or non-101 status code.', response.error) return } // HTTP2 if (response.status !== 200) { - failWebsocketConnection(handler, 1002, 'Received network error or non-200 status code.', response.error) + failHandshake(handler, response, 1002, 'Received network error or non-200 status code.', response.error) return } } @@ -65255,7 +65857,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // header list results in null, failure, or the empty byte // sequence, then fail the WebSocket connection. if (protocols.length !== 0 && !response.headersList.get('Sec-WebSocket-Protocol')) { - failWebsocketConnection(handler, 1002, 'Server did not respond with sent protocols.') + failHandshake(handler, response, 1002, 'Server did not respond with sent protocols.') return } @@ -65271,7 +65873,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // _Fail the WebSocket Connection_. // For H2, no upgrade header is expected. if (response.socket.session == null && response.headersList.get('Upgrade')?.toLowerCase() !== 'websocket') { - failWebsocketConnection(handler, 1002, 'Server did not set Upgrade header to "websocket".') + failHandshake(handler, response, 1002, 'Server did not set Upgrade header to "websocket".') return } @@ -65281,7 +65883,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // MUST _Fail the WebSocket Connection_. // For H2, no connection header is expected. if (response.socket.session == null && response.headersList.get('Connection')?.toLowerCase() !== 'upgrade') { - failWebsocketConnection(handler, 1002, 'Server did not set Connection header to "upgrade".') + failHandshake(handler, response, 1002, 'Server did not set Connection header to "upgrade".') return } @@ -65295,7 +65897,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) const secWSAccept = response.headersList.get('Sec-WebSocket-Accept') const digest = crypto.hash('sha1', keyValue + uid, 'base64') if (secWSAccept !== digest) { - failWebsocketConnection(handler, 1002, 'Incorrect hash received in Sec-WebSocket-Accept header.') + failHandshake(handler, response, 1002, 'Incorrect hash received in Sec-WebSocket-Accept header.') return } @@ -65313,7 +65915,7 @@ function establishWebSocketConnection (url, protocols, client, handler, options) extensions = parseExtensions(secExtension) if (!extensions.has('permessage-deflate')) { - failWebsocketConnection(handler, 1002, 'Sec-WebSocket-Extensions header does not match.') + failHandshake(handler, response, 1002, 'Sec-WebSocket-Extensions header does not match.') return } } @@ -65333,8 +65935,8 @@ function establishWebSocketConnection (url, protocols, client, handler, options) // 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)) { - failWebsocketConnection(handler, 1002, 'Protocol was not set in the opening handshake.') + if (requestProtocols === null || !requestProtocols.includes(secProtocol)) { + failHandshake(handler, response, 1002, 'Protocol was not set in the opening handshake.') return } } @@ -65429,6 +66031,16 @@ function closeWebSocketConnection (object, code, reason, validate = false) { } } +function failHandshake (handler, response, code, reason, cause) { + // The H2 upgrade request has already completed and handed off its stream. + // Aborting the request cannot close that stream after handshake validation fails. + if (response.socket?.session != null && !response.socket.destroyed) { + response.socket.destroy() + } + + failWebsocketConnection(handler, code, reason, cause) +} + /** * @param {import('./websocket').Handler} handler * @param {number} code @@ -66147,7 +66759,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 } @@ -66970,9 +67587,9 @@ class WebSocketStream { /** @type {ReadableStreamDefaultController} */ #readableStreamController - // Each WebSocketStream object has an associated writable stream , which is a WritableStream . - /** @type {WritableStream} */ - #writableStream + // Retain the controller so the writable stream can be errored while locked. + /** @type {WritableStreamDefaultController} */ + #writableStreamController // Each WebSocketStream object has an associated boolean handshake aborted , which is initially false. #handshakeAborted = false @@ -67236,6 +67853,9 @@ class WebSocketStream { // 12. Let writable be a new WritableStream . // 13. Set up writable with writeAlgorithm , closeAlgorithm , and abortAlgorithm . const writable = new WritableStream({ + start: (controller) => { + this.#writableStreamController = controller + }, write: (chunk) => this.#write(chunk), close: () => closeWebSocketConnection(this.#handler, null, null), abort: (reason) => this.#closeUsingReason(reason) @@ -67244,9 +67864,6 @@ class WebSocketStream { // Set stream ’s readable stream to readable . this.#readableStream = readable - // Set stream ’s writable stream to writable . - this.#writableStream = writable - // Resolve stream ’s opened promise with WebSocketOpenInfo «[ " extensions " → extensions , " protocol " → protocol , " readable " → readable , " writable " → writable ]». this.#openedPromise.resolve({ extensions, @@ -67332,9 +67949,7 @@ class WebSocketStream { this.#readableStreamController.close() // 6.2. Error stream ’s writable stream with an " InvalidStateError " DOMException indicating that a closed WebSocketStream cannot be written to. - if (!this.#writableStream.locked) { - this.#writableStream.abort(new DOMException('A closed WebSocketStream cannot be written to', 'InvalidStateError')) - } + this.#writableStreamController.error(new DOMException('A closed WebSocketStream cannot be written to', 'InvalidStateError')) // 6.3. Resolve stream ’s closed promise with WebSocketCloseInfo «[ " closeCode " → code , " reason " → reason ]». this.#closedPromise.resolve({ @@ -67351,7 +67966,7 @@ class WebSocketStream { this.#readableStreamController?.error(error) // 7.3. Error stream ’s writable stream with error . - this.#writableStream?.abort(error) + this.#writableStreamController?.error(error) // 7.4. Reject stream ’s closed promise with error . this.#closedPromise.reject(error) diff --git a/package-lock.json b/package-lock.json index 22b2ed6..5314a58 100644 --- a/package-lock.json +++ b/package-lock.json @@ -4233,9 +4233,9 @@ } }, "node_modules/undici": { - "version": "7.29.0", - "resolved": "https://registry.npmjs.org/undici/-/undici-7.29.0.tgz", - "integrity": "sha512-IDxfleLmmbSskfWSUATiN1nfn2rDuvnMOqb5CWR92iIfojA0Ud+ulOAAEQ57LPr9rWmsreUyf5lwyao+7GNNVw==", + "version": "7.30.0", + "resolved": "https://registry.npmjs.org/undici/-/undici-7.30.0.tgz", + "integrity": "sha512-dkrQXeHSaoamnItlYbmzG0wFYrM0ZwDxCIg0A7aKjTyyhh9svRzCNFEzV+Vm05/yehjCzjDZ31KXfGEjYSztDQ==", "license": "MIT", "engines": { "node": ">=20.18.1"