| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent cc19107 commit f45dc92
21 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -111,22 +111,35 @@ added: v1.0.0 | |||
| 111 | 111 | `autoSelectFamily` is enabled. **Default:** `250`. | |
| 112 | 112 | * `allowH2` {boolean} Enables HTTP/2 support when the server assigns it a | |
| 113 | 113 | higher priority through ALPN negotiation. **Default:** `true`. | |
| 114 | - * `useH2c` {boolean} Enforces h2c (HTTP/2 cleartext) for non-HTTPS | ||
| 115 | - connections. **Default:** `false`. | ||
| 116 | - * `maxConcurrentStreams` {number} The maximum number of concurrent HTTP/2 | ||
| 114 | + * `useH2c` {boolean} _Deprecated: use h2Options.useH2c instead_ Enforces h2c (HTTP/2 cleartext) for non-HTTPS | ||
| 115 | + connections. **Default:** `false`. | ||
| 116 | + * `maxConcurrentStreams` {number} _Deprecated: use h2Options.useH2c instead_ The maximum number of concurrent HTTP/2 | ||
| 117 | 117 | streams for a single session. Once h2 is negotiated this — not `pipelining`, | |
| 118 | 118 | which is HTTP/1.1 only — is the ceiling used to dispatch in-flight requests. | |
| 119 | 119 | It may be overridden by the server's `SETTINGS_MAX_CONCURRENT_STREAMS` | |
| 120 | 120 | frame. **Default:** `100`. | |
| 121 | - * `initialWindowSize` {number} The HTTP/2 stream-level flow-control window | ||
| 122 | - size (`SETTINGS_INITIAL_WINDOW_SIZE`). Must be a positive integer. | ||
| 123 | - **Default:** `262144`. | ||
| 124 | - * `connectionWindowSize` {number} The HTTP/2 connection-level flow-control | ||
| 121 | + * `connectionWindowSize` {number} _Deprecated: use h2Options.connectionWindowSize instead_ The HTTP/2 connection-level flow-control | ||
| 125 | 122 | window size set via `ClientHttp2Session.setLocalWindowSize()`. Must be a | |
| 126 | 123 | positive integer. **Default:** `524288`. | |
| 127 | - * `pingInterval` {number} The time interval, in milliseconds, between HTTP/2 | ||
| 124 | + * `pingInterval` {number} _Deprecated: use h2Options.pingInterval instead_ The time interval, in milliseconds, between HTTP/2 | ||
| 128 | 125 | PING frames. Set to `0` to disable PING frames. Applies only to HTTP/2 | |
| 129 | 126 | connections and emits a `ping` event on the client. **Default:** `60e3`. | |
| 127 | + * `h2Options` {object} Set of options for HTTP/2 sessions | ||
| 128 | + * `useH2c` {boolean} Enforces h2c (HTTP/2 cleartext) for non-HTTPS | ||
| 129 | + connections. **Default:** `false`. | ||
| 130 | + * `maxConcurrentStreams` {number} The maximum number of concurrent HTTP/2 | ||
| 131 | + streams for a single session. Once h2 is negotiated this — not `pipelining`, | ||
| 132 | + which is HTTP/1.1 only — is the ceiling used to dispatch in-flight requests. | ||
| 133 | + It may be overridden by the server's `SETTINGS_MAX_CONCURRENT_STREAMS` | ||
| 134 | + frame. **Default:** `100`. | ||
| 135 | + * `connectionWindowSize` {number} The HTTP/2 connection-level flow-control | ||
| 136 | + window size set via `ClientHttp2Session.setLocalWindowSize()`. Must be a | ||
| 137 | + positive integer. **Default:** `524288`. | ||
| 138 | + * `pingInterval` {number} The time interval, in milliseconds, between HTTP/2 | ||
| 139 | + PING frames. Set to `0` to disable PING frames. Applies only to HTTP/2 | ||
| 140 | + connections and emits a `ping` event on the client. **Default:** `60e3`. | ||
| 141 | + * `settings` {object} `SETTINGS` frame options. For full reference, take a | ||
| 142 | + look to [HTTP/2#Settings Object](https://nodejs.org/api/http2.html#settings-object) | ||
| 130 | 143 | * `webSocket` {Object} (optional) WebSocket-specific configuration. | |
| 131 | 144 | * `maxFragments` {number} The maximum number of fragments in a message. Set | |
| 132 | 145 | to `0` to disable the limit. **Default:** `131072`. | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -15,7 +15,6 @@ const kContentType = Symbol('kContentType') | |||
| 15 | 15 | const kContentLength = Symbol('kContentLength') | |
| 16 | 16 | const kUsed = Symbol('kUsed') | |
| 17 | 17 | const kBytesRead = Symbol('kBytesRead') | |
| 18 | - const kPreservedBuffer = Symbol('kPreservedBuffer') | ||
| 19 | 18 | ||
| 20 | 19 | const noop = () => {} | |
| 21 | 20 | ||
@@ -326,36 +325,14 @@ class BodyReadable extends Readable { | |||
| 326 | 325 | */ | |
| 327 | 326 | setEncoding (encoding) { | |
| 328 | 327 | if (Buffer.isEncoding(encoding)) { | |
| 329 | - // Preserve raw Buffer chunks for the consume path (body.text(), | ||
| 330 | - // body.json(), etc.) before super.setEncoding() replaces them | ||
| 331 | - // with decoded strings. Without this, the consume path would | ||
| 332 | - // lose access to the original bytes — some of which may be held | ||
| 333 | - // by the decoder for incomplete multi-byte sequences, and the | ||
| 334 | - // rest converted to strings that can't be safely concatenated | ||
| 335 | - // byte-wise. | ||
| 336 | - const state = this._readableState | ||
| 337 | - const buffer = state.buffer | ||
| 338 | - if (buffer && state.length > 0) { | ||
| 339 | - const bufferIndex = state.bufferIndex ?? 0 | ||
| 340 | - const preserved = [] | ||
| 341 | - const source = typeof buffer.slice === 'function' | ||
| 342 | - ? buffer.slice(bufferIndex) | ||
| 343 | - : buffer | ||
| 344 | - for (const data of source) { | ||
| 345 | - if (Buffer.isBuffer(data)) { | ||
| 346 | - preserved.push(data) | ||
| 347 | - } | ||
| 348 | - } | ||
| 349 | - if (preserved.length > 0) { | ||
| 350 | - this[kPreservedBuffer] = (this[kPreservedBuffer] || []).concat(preserved) | ||
| 351 | - } | ||
| 352 | - } | ||
| 353 | - | ||
| 354 | 328 | // Delegate to Node.js Readable.setEncoding() which initializes a | |
| 355 | 329 | // StringDecoder and re-encodes already-buffered chunks. This properly | |
| 356 | 330 | // handles multi-byte sequences split at chunk boundaries for the | |
| 357 | 331 | // for-await / on('data') paths. Without this, Node.js uses | |
| 358 | 332 | // buf.toString(encoding) on each chunk, producing U+FFFD for split chars. | |
| 333 | + // | ||
| 334 | + // The consume path (body.text(), body.json(), ...) copes with the | ||
| 335 | + // decoded strings this leaves in state.buffer, see consumeStart(). | ||
| 359 | 336 | super.setEncoding(encoding) | |
| 360 | 337 | } | |
| 361 | 338 | return this | |
@@ -464,17 +441,7 @@ function consumeStart (consume) { | |||
| 464 | 441 | ||
| 465 | 442 | const { _readableState: state } = consume.stream | |
| 466 | 443 | ||
| 467 | - // If setEncoding() was called, state.buffer may contain decoded strings | ||
| 468 | - // (which would break Buffer.concat in chunksDecode). Use the preserved | ||
| 469 | - // raw Buffers (saved before super.setEncoding() in setEncoding()) for | ||
| 470 | - // byte-level accurate consumption. Otherwise read from state.buffer. | ||
| 471 | - const preserved = consume.stream[kPreservedBuffer] | ||
| 472 | - if (preserved && preserved.length > 0) { | ||
| 473 | - for (const chunk of preserved) { | ||
| 474 | - consumePush(consume, chunk) | ||
| 475 | - } | ||
| 476 | - consume.stream[kPreservedBuffer] = null | ||
| 477 | - } else if (state.bufferIndex) { | ||
| 444 | + if (state.bufferIndex) { | ||
| 478 | 445 | const start = state.bufferIndex | |
| 479 | 446 | const end = state.buffer.length | |
| 480 | 447 | for (let n = start; n < end; n++) { | |
@@ -486,14 +453,29 @@ function consumeStart (consume) { | |||
| 486 | 453 | } | |
| 487 | 454 | } | |
| 488 | 455 | ||
| 456 | + // If setEncoding() was called, state.buffer holds decoded strings, which | ||
| 457 | + // consumePush() turns back into bytes. The trailing bytes of a multi-byte | ||
| 458 | + // sequence split across a chunk boundary are not part of any of those | ||
| 459 | + // strings, they are held inside the decoder until the rest arrives, so | ||
| 460 | + // take them from there. | ||
| 461 | + const decoder = state.decoder | ||
| 462 | + if (decoder != null && decoder.lastNeed > 0) { | ||
| 463 | + consumePush(consume, Buffer.from(decoder.lastChar.subarray(0, decoder.lastTotal - decoder.lastNeed))) | ||
| 464 | + } | ||
| 465 | + | ||
| 489 | 466 | if (state.endEmitted) { | |
| 490 | - consumeEnd(this[kConsume], this._readableState.encoding) | ||
| 491 | - } else { | ||
| 492 | - consume.stream.on('end', function () { | ||
| 493 | - consumeEnd(this[kConsume], this._readableState.encoding) | ||
| 494 | - }) | ||
| 467 | + // No `this` to read the consume off here: consumeStart is a free function, called from | ||
| 468 | + // the queueMicrotask above. The callback below does have one, because the emitter passes | ||
| 469 | + // the stream as its receiver. Returning matters too - consumeEnd() clears consume.stream, | ||
| 470 | + // which the resume() below would then dereference. | ||
| 471 | + consumeEnd(consume, state.encoding) | ||
| 472 | + return | ||
| 495 | 473 | } | |
| 496 | 474 | ||
| 475 | + consume.stream.on('end', function () { | ||
| 476 | + consumeEnd(this[kConsume], this._readableState.encoding) | ||
| 477 | + }) | ||
| 478 | + | ||
| 497 | 479 | consume.stream.resume() | |
| 498 | 480 | ||
| 499 | 481 | while (consume.stream.read() != null) { | |
@@ -583,14 +565,22 @@ function consumeEnd (consume, encoding) { | |||
| 583 | 565 | ||
| 584 | 566 | /** | |
| 585 | 567 | * @param {Consume} consume | |
| 586 | - * @param {Buffer} chunk | ||
| 568 | + * @param {Buffer|string} chunk | ||
| 587 | 569 | * @returns {void} | |
| 588 | 570 | */ | |
| 589 | 571 | function consumePush (consume, chunk) { | |
| 590 | 572 | if (consume.body === null) { | |
| 591 | 573 | return | |
| 592 | 574 | } | |
| 593 | 575 | ||
| 576 | + if (typeof chunk === 'string') { | ||
| 577 | + // Buffered before the consume started, while an encoding was set. | ||
| 578 | + // consume.length has to stay a byte count and chunksDecode()/chunksConcat() | ||
| 579 | + // only work on bytes, so re-encode. A string's own length is in UTF-16 code | ||
| 580 | + // units and Uint8Array.prototype.set() ignores a string argument entirely. | ||
| 581 | + chunk = Buffer.from(chunk, consume.stream._readableState.encoding) | ||
| 582 | + } | ||
| 583 | + | ||
| 594 | 584 | consume.length += chunk.length | |
| 595 | 585 | consume.body.push(chunk) | |
| 596 | 586 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -105,13 +105,27 @@ function buildConnector ({ allowH2, preferH2, useH2c, maxCachedSessions, socketP | |||
| 105 | 105 | ||
| 106 | 106 | port = port || 80 | |
| 107 | 107 | ||
| 108 | - socket = net.connect({ | ||
| 108 | + const connectOptions = { | ||
| 109 | 109 | highWaterMark: 64 * 1024, // Same as nodejs fs streams. | |
| 110 | 110 | ...options, | |
| 111 | 111 | localAddress, | |
| 112 | 112 | port, | |
| 113 | 113 | host: hostname | |
| 114 | - }) | ||
| 114 | + } | ||
| 115 | + | ||
| 116 | + const family = net.isIP(hostname) | ||
| 117 | + if (family !== 0 && servername && servername !== hostname) { | ||
| 118 | + connectOptions.host = servername | ||
| 119 | + connectOptions.lookup = (_hostname, lookupOptions, cb) => { | ||
| 120 | + if (lookupOptions.all) { | ||
| 121 | + cb(null, [{ address: hostname, family }]) | ||
| 122 | + } else { | ||
| 123 | + cb(null, hostname, family) | ||
| 124 | + } | ||
| 125 | + } | ||
| 126 | + } | ||
| 127 | + | ||
| 128 | + socket = net.connect(connectOptions) | ||
| 115 | 129 | if (useH2c === true) { | |
| 116 | 130 | socket.alpnProtocol = 'h2' | |
| 117 | 131 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -56,6 +56,7 @@ module.exports = { | |||
| 56 | 56 | kCounter: Symbol('socket request counter'), | |
| 57 | 57 | kMaxResponseSize: Symbol('max response size'), | |
| 58 | 58 | kHTTP2Session: Symbol('http2Session'), | |
| 59 | + kHTTP2Options: Symbol('http2 options'), | ||
| 59 | 60 | kHTTP2SessionState: Symbol('http2Session state'), | |
| 60 | 61 | kRetryHandlerDefaultRetry: Symbol('retry agent default retry'), | |
| 61 | 62 | kConstruct: Symbol('constructable'), | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1052,7 +1052,7 @@ function onSocketClose () { | |||
| 1052 | 1052 | ||
| 1053 | 1053 | function clearIdleSocketValidation (socket) { | |
| 1054 | 1054 | if (socket[kIdleSocketValidationTimeout]) { | |
| 1055 | - clearImmediate(socket[kIdleSocketValidationTimeout]) | ||
| 1055 | + clearTimeout(socket[kIdleSocketValidationTimeout]) | ||
| 1056 | 1056 | socket[kIdleSocketValidationTimeout] = null | |
| 1057 | 1057 | } | |
| 1058 | 1058 | ||
@@ -1061,14 +1061,14 @@ function clearIdleSocketValidation (socket) { | |||
| 1061 | 1061 | ||
| 1062 | 1062 | function scheduleIdleSocketValidation (client, socket) { | |
| 1063 | 1063 | socket[kIdleSocketValidation] = 1 | |
| 1064 | - socket[kIdleSocketValidationTimeout] = setImmediate(() => { | ||
| 1064 | + socket[kIdleSocketValidationTimeout] = setTimeout(() => { | ||
| 1065 | 1065 | socket[kIdleSocketValidationTimeout] = null | |
| 1066 | 1066 | socket[kIdleSocketValidation] = 2 | |
| 1067 | 1067 | ||
| 1068 | 1068 | if (client[kSocket] === socket && !socket.destroyed) { | |
| 1069 | 1069 | client[kResume]() | |
| 1070 | 1070 | } | |
| 1071 | - }) | ||
| 1071 | + }, 0) | ||
| 1072 | 1072 | socket[kIdleSocketValidationTimeout].unref?.() | |
| 1073 | 1073 | } | |
| 1074 | 1074 | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments