diff --git a/lib/entry-points.js b/lib/entry-points.js index d106883fa..54d58abd9 100644 --- a/lib/entry-points.js +++ b/lib/entry-points.js @@ -2148,10 +2148,66 @@ var require_request = __commonJS({ return false; } } - onUpgrade(statusCode, headers, socket) { + /** + * @param {number|null} statusCode + * @param {Buffer[]|null} headers + * @param {import('node:stream').Duplex} socket + * @param {string} [statusText] + */ + onUpgrade(statusCode, headers, socket, statusText = "") { + this.onFinally(); assert(!this.aborted); assert(!this.completed); - return this[kHandler].onUpgrade(statusCode, headers, socket); + if (statusCode !== null) { + this.#publishUpgradeHeaders(statusCode, headers, statusText); + } + const result = this[kHandler].onUpgrade(statusCode, headers, socket); + if (!this.aborted) { + this.completed = true; + if (statusCode !== null) { + this.#publishUpgradeTrailers(); + } + } + return result; + } + /** + * @param {number} statusCode + * @param {import('node:http2').IncomingHttpHeaders} headers + * @param {(headers: import('node:http2').IncomingHttpHeaders) => Buffer[]} parseHeaders + * @param {string} [statusText] + */ + onUpgradeResponse(statusCode, headers, parseHeaders, statusText = "") { + assert(!this.aborted); + assert(this.completed); + if (channels.headers.hasSubscribers) { + this.#publishUpgradeHeaders(statusCode, parseHeaders(headers), statusText); + } + this.#publishUpgradeTrailers(); + } + /** + * @param {Error} error + */ + onUpgradeError(error3) { + assert(!this.aborted); + assert(this.completed); + if (channels.error.hasSubscribers) { + channels.error.publish({ request: this, error: error3 }); + } + } + /** + * @param {number} statusCode + * @param {Buffer[]} headers + * @param {string} statusText + */ + #publishUpgradeHeaders(statusCode, headers, statusText) { + if (channels.headers.hasSubscribers) { + channels.headers.publish({ request: this, response: { statusCode, headers, statusText } }); + } + } + #publishUpgradeTrailers() { + if (channels.trailers.hasSubscribers) { + channels.trailers.publish({ request: this, trailers: [] }); + } } onComplete(trailers) { this.onFinally(); @@ -6294,7 +6350,7 @@ var require_client_h1 = __commonJS({ } } onUpgrade(head) { - const { upgrade, client, socket, headers, statusCode } = this; + const { upgrade, client, socket, headers, statusCode, statusText } = this; assert(upgrade); assert(client[kSocket] === socket); assert(!socket.destroyed); @@ -6319,9 +6375,10 @@ var require_client_h1 = __commonJS({ client[kQueue][client[kRunningIdx]++] = null; client.emit("disconnect", client[kUrl], [client], new InformationalError("upgrade")); try { - request3.onUpgrade(statusCode, headers, socket); - } catch (err) { - util4.destroy(socket, err); + request3.onUpgrade(statusCode, headers, socket, statusText); + } catch (error3) { + util4.errorRequest(client, request3, error3); + util4.destroy(socket, error3); } client[kResume](); } @@ -6613,21 +6670,20 @@ var require_client_h1 = __commonJS({ } function clearIdleSocketValidation(socket) { if (socket[kIdleSocketValidationTimeout]) { - clearTimeout(socket[kIdleSocketValidationTimeout]); + clearImmediate(socket[kIdleSocketValidationTimeout]); socket[kIdleSocketValidationTimeout] = null; } socket[kIdleSocketValidation] = 0; } function scheduleIdleSocketValidation(client, socket) { socket[kIdleSocketValidation] = 1; - socket[kIdleSocketValidationTimeout] = setTimeout(() => { + socket[kIdleSocketValidationTimeout] = setImmediate(() => { socket[kIdleSocketValidationTimeout] = null; socket[kIdleSocketValidation] = 2; if (client[kSocket] === socket && !socket.destroyed) { client[kResume](); } - }, 0); - socket[kIdleSocketValidationTimeout].unref?.(); + }); } function resumeH1(client) { const socket = client[kSocket]; @@ -6725,11 +6781,17 @@ var require_client_h1 = __commonJS({ } const socket = client[kSocket]; clearIdleSocketValidation(socket); - const abort = (err) => { - if (request3.aborted || request3.completed) { + const abort = (error3) => { + if (request3.aborted) { return; } - util4.errorRequest(client, request3, err || new RequestAbortedError()); + if (request3.completed) { + if (request3.upgrade || request3.method === "CONNECT") { + util4.destroy(socket, new InformationalError("aborted")); + } + return; + } + util4.errorRequest(client, request3, error3 || new RequestAbortedError()); util4.destroy(body); util4.destroy(socket, new InformationalError("aborted")); }; @@ -7085,6 +7147,7 @@ var require_client_h2 = __commonJS({ "node_modules/undici/lib/dispatcher/client-h2.js"(exports2, module2) { "use strict"; var assert = require("node:assert"); + var { errorMonitor } = require("node:events"); var { pipeline: pipeline2 } = require("node:stream"); var util4 = require_util(); var { @@ -7145,6 +7208,10 @@ var require_client_h2 = __commonJS({ } return result; } + function parseH2ResponseHeaders(headers) { + const { [HTTP2_HEADER_STATUS]: _statusCode, ...realHeaders } = headers; + return parseH2Headers(realHeaders); + } async function connectH2(client, socket) { client[kSocket] = socket; if (!h2ExperimentalWarned) { @@ -7308,16 +7375,22 @@ var require_client_h2 = __commonJS({ const { hostname, port } = client[kUrl]; headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ""}`; headers[HTTP2_HEADER_METHOD] = method; - const abort = (err) => { - if (request3.aborted || request3.completed) { + const abort = (error3) => { + if (request3.aborted) { return; } - err = err || new RequestAbortedError(); - util4.errorRequest(client, request3, err); - if (stream2 != null) { - util4.destroy(stream2, err); + if (request3.completed) { + if (method === "CONNECT" && stream2 != null) { + util4.destroy(stream2, error3 || new RequestAbortedError()); + } + return; } - util4.destroy(body, err); + error3 = error3 || new RequestAbortedError(); + util4.errorRequest(client, request3, error3); + if (stream2 != null) { + util4.destroy(stream2, error3); + } + util4.destroy(body, error3); client[kQueue][client[kRunningIdx]++] = null; client[kResume](); }; @@ -7332,18 +7405,42 @@ var require_client_h2 = __commonJS({ if (method === "CONNECT") { session.ref(); stream2 = session.request(headers, { endStream: false, signal }); - if (stream2.id && !stream2.pending) { - request3.onUpgrade(null, null, stream2); - ++session[kOpenStreams]; - client[kQueue][client[kRunningIdx]++] = null; - } else { - stream2.once("ready", () => { + let upgradeResponseFinished = false; + const onResponse = (headers2) => { + upgradeResponseFinished = true; + stream2.off(errorMonitor, onUpgradeError); + request3.onUpgradeResponse(Number(headers2[HTTP2_HEADER_STATUS]), headers2, parseH2ResponseHeaders); + }; + const onUpgradeError = (error3) => { + upgradeResponseFinished = true; + stream2.off("response", onResponse); + request3.onUpgradeError(error3); + }; + const onReady = () => { + try { request3.onUpgrade(null, null, stream2); - ++session[kOpenStreams]; - client[kQueue][client[kRunningIdx]++] = null; - }); - } + } catch (error3) { + stream2.off("response", onResponse); + abort(error3); + return; + } + if (request3.aborted) { + return; + } + stream2.off("error", abort); + stream2.once(errorMonitor, onUpgradeError); + client[kQueue][client[kRunningIdx]++] = null; + }; + stream2.once("response", onResponse); + stream2.once("error", abort); + ++session[kOpenStreams]; + onReady(); stream2.once("close", () => { + if (!upgradeResponseFinished && request3.completed) { + stream2.off("response", onResponse); + stream2.off(errorMonitor, onUpgradeError); + request3.onUpgradeError(new InformationalError(`HTTP/2: "stream error" received - code ${stream2.rstCode}`)); + } session[kOpenStreams] -= 1; if (session[kOpenStreams] === 0) session.unref(); }); @@ -9326,6 +9423,7 @@ var require_retry_handler = __commonJS({ this.end = null; this.etag = null; this.resume = null; + this.headersSent = false; this.handler.onConnect((reason) => { this.aborted = true; if (this.abort) { @@ -9335,6 +9433,17 @@ var require_retry_handler = __commonJS({ } }); } + checkpointResponseEnd(headers, resume) { + if (this.end == null && this.opts.method !== "HEAD") { + const contentLength = headers["content-length"]; + this.end = contentLength != null ? Number(contentLength) - 1 : null; + assert( + this.end == null || Number.isFinite(this.end), + "invalid content-length" + ); + } + this.resume = this.end != null ? resume : null; + } onRequestSent() { if (this.handler.onRequestSent) { this.handler.onRequestSent(); @@ -9396,7 +9505,9 @@ var require_retry_handler = __commonJS({ const headers = parseHeaders(rawHeaders); this.retryCount += 1; if (statusCode >= 300) { - if (this.retryOpts.statusCodes.includes(statusCode) === false) { + if (!this.headersSent && this.retryOpts.statusCodes.includes(statusCode) === false) { + this.headersSent = true; + this.checkpointResponseEnd(headers, resume); return this.handler.onHeaders( statusCode, rawHeaders, @@ -9451,8 +9562,15 @@ var require_retry_handler = __commonJS({ return false; } const { start, size, end = size - 1 } = contentRange; - assert(this.start === start, "content-range mismatch"); - assert(this.end == null || this.end === end, "content-range mismatch"); + if (this.start !== start || this.end != null && this.end !== end) { + this.abort( + new RequestRetryError("Content-Range mismatch", statusCode, { + headers, + data: { count: this.retryCount } + }) + ); + return false; + } this.resume = resume; return true; } @@ -9460,6 +9578,7 @@ var require_retry_handler = __commonJS({ if (statusCode === 206) { const range2 = parseRangeHeader(headers["content-range"]); if (range2 == null) { + this.headersSent = true; return this.handler.onHeaders( statusCode, rawHeaders, @@ -9491,6 +9610,7 @@ var require_retry_handler = __commonJS({ "invalid content-length" ); this.resume = resume; + this.headersSent = true; this.etag = headers.etag != null ? headers.etag : null; if (this.etag != null && this.etag.startsWith("W/")) { this.etag = null; @@ -9518,7 +9638,7 @@ var require_retry_handler = __commonJS({ return this.handler.onComplete(rawTrailers); } onError(err) { - if (this.aborted || isDisturbed(this.opts.body)) { + if (this.aborted || isDisturbed(this.opts.body) || this.headersSent && this.resume == null) { return this.handler.onError(err); } if (this.retryCount - this.retryCountCheckpoint > 0) { @@ -17404,7 +17524,7 @@ var require_connection = __commonJS({ const secProtocol = response.headersList.get("Sec-WebSocket-Protocol"); if (secProtocol !== null) { const requestProtocols = getDecodeSplit("sec-websocket-protocol", request3.headersList); - if (!requestProtocols.includes(secProtocol)) { + if (requestProtocols === null || !requestProtocols.includes(secProtocol)) { failWebsocketConnection(ws, "Protocol was not set in the opening handshake."); return; } @@ -17552,6 +17672,7 @@ var require_permessage_deflate = __commonJS({ if (this.#maxPayloadSize > 0 && this.#inflate[kLength] > this.#maxPayloadSize) { callback(new MessageSizeExceededError()); this.#inflate.removeAllListeners(); + this.#inflate.destroy(); this.#inflate = null; return; } @@ -18463,6 +18584,40 @@ var require_eventsource_stream = __commonJS({ var CR = 13; var COLON = 58; var SPACE = 32; + var DATA = Buffer.from("data"); + var EVENT = Buffer.from("event"); + var ID2 = Buffer.from("id"); + var 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] < 48 || buffer[i] > 57) { + return false; + } + } + return true; + } + function isValidLastEventIdBytes(buffer, start) { + for (let i = start; i < buffer.length; i++) { + if (buffer[i] === 0) { + 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; + } var EventSourceStream = class extends Transform5 { /** * @type {eventSourceSettings} @@ -18482,10 +18637,13 @@ var require_eventsource_stream = __commonJS({ */ eventEndCheck = false; /** - * @type {Buffer} + * @type {Buffer[]} */ - buffer = null; + chunks = []; + chunkIndex = 0; pos = 0; + lineChunkIndex = 0; + linePos = 0; event = { data: void 0, event: void 0, @@ -18516,63 +18674,30 @@ var require_eventsource_stream = __commonJS({ callback(); return; } - if (this.buffer) { - this.buffer = Buffer.concat([this.buffer, chunk]); - } else { - this.buffer = chunk; - } + this.chunks.push(chunk); if (this.checkBOM) { - switch (this.buffer.length) { - case 1: - if (this.buffer[0] === BOM[0]) { - callback(); - return; - } - this.checkBOM = false; - callback(); - return; - case 2: - if (this.buffer[0] === BOM[0] && this.buffer[1] === BOM[1]) { - callback(); - return; - } - this.checkBOM = false; - break; - case 3: - if (this.buffer[0] === BOM[0] && this.buffer[1] === BOM[1] && this.buffer[2] === BOM[2]) { - this.buffer = Buffer.alloc(0); - this.checkBOM = false; - callback(); - return; - } - this.checkBOM = false; - break; - default: - if (this.buffer[0] === BOM[0] && this.buffer[1] === BOM[1] && this.buffer[2] === BOM[2]) { - this.buffer = this.buffer.subarray(3); - } - this.checkBOM = false; - break; + if (this.handleBOM()) { + callback(); + return; } } - while (this.pos < this.buffer.length) { + while (this.hasCurrentByte()) { + const byte = this.currentByte(); if (this.eventEndCheck) { if (this.crlfCheck) { - 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(); continue; } this.crlfCheck = false; } - if (this.buffer[this.pos] === LF || this.buffer[this.pos] === CR) { - if (this.buffer[this.pos] === CR) { + if (byte === LF || byte === CR) { + if (byte === CR) { this.crlfCheck = true; } - this.buffer = this.buffer.subarray(this.pos + 1); - this.pos = 0; - if (this.event.data !== void 0 || this.event.event || this.event.id || this.event.retry) { + this.consumeCurrentByte(); + if (this.hasPendingEvent()) { this.processEvent(this.event); } this.clearEvent(); @@ -18581,17 +18706,16 @@ var require_eventsource_stream = __commonJS({ this.eventEndCheck = false; continue; } - if (this.buffer[this.pos] === LF || this.buffer[this.pos] === CR) { - if (this.buffer[this.pos] === CR) { + if (byte === LF || byte === CR) { + if (byte === CR) { this.crlfCheck = true; } - this.parseLine(this.buffer.subarray(0, this.pos), this.event); - this.buffer = this.buffer.subarray(this.pos + 1); - this.pos = 0; + this.parseLine(this.readLine(), this.event); + this.consumeCurrentByte(); this.eventEndCheck = true; continue; } - this.pos++; + this.advanceCursor(); } callback(); } @@ -18607,43 +18731,42 @@ var require_eventsource_stream = __commonJS({ if (colonPosition === 0) { return; } - let field = ""; - let value = ""; + let fieldLength = line.length; + let valueStart = line.length; if (colonPosition !== -1) { - field = line.subarray(0, colonPosition).toString("utf8"); - let valueStart = colonPosition + 1; + fieldLength = colonPosition; + valueStart = colonPosition + 1; if (line[valueStart] === SPACE) { ++valueStart; } - value = line.subarray(valueStart).toString("utf8"); - } else { - field = line.toString("utf8"); - value = ""; } - switch (field) { - case "data": - if (event[field] === void 0) { - event[field] = value; - } else { - event[field] += ` + if (isFieldName(line, fieldLength, DATA)) { + const value = line.toString("utf8", valueStart); + if (event.data === void 0) { + event.data = value; + } else { + event.data += ` ${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; + } + return; + } + if (isFieldName(line, fieldLength, RETRY)) { + if (isASCIINumberBytes(line, valueStart)) { + event.retry = line.toString("utf8", valueStart); + } + return; + } + if (isFieldName(line, fieldLength, ID2)) { + 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; + } } } /** @@ -18668,12 +18791,120 @@ ${value}`; } } clearEvent() { - this.event = { - data: void 0, - event: void 0, - id: void 0, - retry: void 0 - }; + this.event.data = void 0; + this.event.event = void 0; + this.event.id = void 0; + this.event.retry = void 0; + } + hasPendingEvent() { + return this.event.data !== void 0 || this.event.event !== void 0 || this.event.id !== void 0 || this.event.retry !== void 0; + } + 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 === void 0) { + if (first === BOM[0]) { + return true; + } + this.checkBOM = false; + return true; + } + if (third === void 0) { + 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(); } }; module2.exports = { diff --git a/package-lock.json b/package-lock.json index cd18608ef..d47461d6c 100644 --- a/package-lock.json +++ b/package-lock.json @@ -36,7 +36,7 @@ "long": "^5.3.2", "node-forge": "^1.4.0", "semver": "^7.8.5", - "undici": "^6.28.0", + "undici": "^6.29.0", "uuid": "^14.0.2" }, "devDependencies": { @@ -10039,9 +10039,9 @@ } }, "node_modules/undici": { - "version": "6.28.0", - "resolved": "https://registry.npmjs.org/undici/-/undici-6.28.0.tgz", - "integrity": "sha512-LIY910g9TI13YS95lrMFrs8Rm/u/irgHeTWoKCoteeJ04CUJ92eEfj0rVn+7VKMPBpUPiUoBKfhNyLI23EE/KA==", + "version": "6.29.0", + "resolved": "https://registry.npmjs.org/undici/-/undici-6.29.0.tgz", + "integrity": "sha512-R+RODBqp6i2pPflGdq+xIOUkl+RNfGgHwoinecKu/JCuf2uO06cOKoDbI2P7Dn6KcswdKwrczbU6IYJ6K8X+wg==", "license": "MIT", "engines": { "node": ">=18.17" diff --git a/package.json b/package.json index 7a4ec0e4e..1bb830f55 100644 --- a/package.json +++ b/package.json @@ -45,7 +45,7 @@ "node-forge": "^1.4.0", "semver": "^7.8.5", "uuid": "^14.0.2", - "undici": "^6.28.0" + "undici": "^6.29.0" }, "devDependencies": { "@ava/typescript": "6.0.0", @@ -95,6 +95,6 @@ "semver": ">=6.3.1" }, "glob": "^13.0.6", - "undici": "^6.28.0" + "undici": "^6.29.0" } }