client.js 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668
  1. 'use strict'
  2. const assert = require('node:assert')
  3. const net = require('node:net')
  4. const http = require('node:http')
  5. const util = require('../core/util.js')
  6. const { ClientStats } = require('../util/stats.js')
  7. const { channels } = require('../core/diagnostics.js')
  8. const Request = require('../core/request.js')
  9. const DispatcherBase = require('./dispatcher-base')
  10. const {
  11. InvalidArgumentError,
  12. InformationalError,
  13. ClientDestroyedError
  14. } = require('../core/errors.js')
  15. const buildConnector = require('../core/connect.js')
  16. const {
  17. kUrl,
  18. kServerName,
  19. kClient,
  20. kBusy,
  21. kConnect,
  22. kResuming,
  23. kRunning,
  24. kPending,
  25. kSize,
  26. kQueue,
  27. kConnected,
  28. kConnecting,
  29. kNeedDrain,
  30. kKeepAliveDefaultTimeout,
  31. kHostHeader,
  32. kPendingIdx,
  33. kRunningIdx,
  34. kError,
  35. kPipelining,
  36. kKeepAliveTimeoutValue,
  37. kMaxHeadersSize,
  38. kKeepAliveMaxTimeout,
  39. kKeepAliveTimeoutThreshold,
  40. kHeadersTimeout,
  41. kBodyTimeout,
  42. kStrictContentLength,
  43. kConnector,
  44. kMaxRequests,
  45. kCounter,
  46. kClose,
  47. kDestroy,
  48. kDispatch,
  49. kLocalAddress,
  50. kMaxResponseSize,
  51. kOnError,
  52. kHTTPContext,
  53. kMaxConcurrentStreams,
  54. kHTTP2InitialWindowSize,
  55. kHTTP2ConnectionWindowSize,
  56. kResume,
  57. kPingInterval
  58. } = require('../core/symbols.js')
  59. const connectH1 = require('./client-h1.js')
  60. const connectH2 = require('./client-h2.js')
  61. const kClosedResolve = Symbol('kClosedResolve')
  62. const getDefaultNodeMaxHeaderSize = http &&
  63. http.maxHeaderSize &&
  64. Number.isInteger(http.maxHeaderSize) &&
  65. http.maxHeaderSize > 0
  66. ? () => http.maxHeaderSize
  67. : () => { throw new InvalidArgumentError('http module not available or http.maxHeaderSize invalid') }
  68. const noop = () => { }
  69. function getPipelining (client) {
  70. return client[kPipelining] ?? client[kHTTPContext]?.defaultPipelining ?? 1
  71. }
  72. /**
  73. * @type {import('../../types/client.js').default}
  74. */
  75. class Client extends DispatcherBase {
  76. /**
  77. *
  78. * @param {string|URL} url
  79. * @param {import('../../types/client.js').Client.Options} options
  80. */
  81. constructor (url, {
  82. maxHeaderSize,
  83. headersTimeout,
  84. socketTimeout,
  85. requestTimeout,
  86. connectTimeout,
  87. bodyTimeout,
  88. idleTimeout,
  89. keepAlive,
  90. keepAliveTimeout,
  91. maxKeepAliveTimeout,
  92. keepAliveMaxTimeout,
  93. keepAliveTimeoutThreshold,
  94. socketPath,
  95. pipelining,
  96. tls,
  97. strictContentLength,
  98. maxCachedSessions,
  99. connect,
  100. maxRequestsPerClient,
  101. localAddress,
  102. maxResponseSize,
  103. autoSelectFamily,
  104. autoSelectFamilyAttemptTimeout,
  105. // h2
  106. maxConcurrentStreams,
  107. allowH2,
  108. useH2c,
  109. initialWindowSize,
  110. connectionWindowSize,
  111. pingInterval,
  112. webSocket
  113. } = {}) {
  114. if (keepAlive !== undefined) {
  115. throw new InvalidArgumentError('unsupported keepAlive, use pipelining=0 instead')
  116. }
  117. if (socketTimeout !== undefined) {
  118. throw new InvalidArgumentError('unsupported socketTimeout, use headersTimeout & bodyTimeout instead')
  119. }
  120. if (requestTimeout !== undefined) {
  121. throw new InvalidArgumentError('unsupported requestTimeout, use headersTimeout & bodyTimeout instead')
  122. }
  123. if (idleTimeout !== undefined) {
  124. throw new InvalidArgumentError('unsupported idleTimeout, use keepAliveTimeout instead')
  125. }
  126. if (maxKeepAliveTimeout !== undefined) {
  127. throw new InvalidArgumentError('unsupported maxKeepAliveTimeout, use keepAliveMaxTimeout instead')
  128. }
  129. if (maxHeaderSize != null) {
  130. if (!Number.isInteger(maxHeaderSize) || maxHeaderSize < 1) {
  131. throw new InvalidArgumentError('invalid maxHeaderSize')
  132. }
  133. } else {
  134. // If maxHeaderSize is not provided, use the default value from the http module
  135. // or if that is not available, throw an error.
  136. maxHeaderSize = getDefaultNodeMaxHeaderSize()
  137. }
  138. if (socketPath != null && typeof socketPath !== 'string') {
  139. throw new InvalidArgumentError('invalid socketPath')
  140. }
  141. if (connectTimeout != null && (!Number.isFinite(connectTimeout) || connectTimeout < 0)) {
  142. throw new InvalidArgumentError('invalid connectTimeout')
  143. }
  144. if (keepAliveTimeout != null && (!Number.isFinite(keepAliveTimeout) || keepAliveTimeout <= 0)) {
  145. throw new InvalidArgumentError('invalid keepAliveTimeout')
  146. }
  147. if (keepAliveMaxTimeout != null && (!Number.isFinite(keepAliveMaxTimeout) || keepAliveMaxTimeout <= 0)) {
  148. throw new InvalidArgumentError('invalid keepAliveMaxTimeout')
  149. }
  150. if (keepAliveTimeoutThreshold != null && !Number.isFinite(keepAliveTimeoutThreshold)) {
  151. throw new InvalidArgumentError('invalid keepAliveTimeoutThreshold')
  152. }
  153. if (headersTimeout != null && (!Number.isInteger(headersTimeout) || headersTimeout < 0)) {
  154. throw new InvalidArgumentError('headersTimeout must be a positive integer or zero')
  155. }
  156. if (bodyTimeout != null && (!Number.isInteger(bodyTimeout) || bodyTimeout < 0)) {
  157. throw new InvalidArgumentError('bodyTimeout must be a positive integer or zero')
  158. }
  159. if (connect != null && typeof connect !== 'function' && typeof connect !== 'object') {
  160. throw new InvalidArgumentError('connect must be a function or an object')
  161. }
  162. if (maxRequestsPerClient != null && (!Number.isInteger(maxRequestsPerClient) || maxRequestsPerClient < 0)) {
  163. throw new InvalidArgumentError('maxRequestsPerClient must be a positive number')
  164. }
  165. if (localAddress != null && (typeof localAddress !== 'string' || net.isIP(localAddress) === 0)) {
  166. throw new InvalidArgumentError('localAddress must be valid string IP address')
  167. }
  168. if (maxResponseSize != null && (!Number.isInteger(maxResponseSize) || maxResponseSize < -1)) {
  169. throw new InvalidArgumentError('maxResponseSize must be a positive number')
  170. }
  171. if (
  172. autoSelectFamilyAttemptTimeout != null &&
  173. (!Number.isInteger(autoSelectFamilyAttemptTimeout) || autoSelectFamilyAttemptTimeout < -1)
  174. ) {
  175. throw new InvalidArgumentError('autoSelectFamilyAttemptTimeout must be a positive number')
  176. }
  177. // h2
  178. if (allowH2 != null && typeof allowH2 !== 'boolean') {
  179. throw new InvalidArgumentError('allowH2 must be a valid boolean value')
  180. }
  181. if (maxConcurrentStreams != null && (typeof maxConcurrentStreams !== 'number' || maxConcurrentStreams < 1)) {
  182. throw new InvalidArgumentError('maxConcurrentStreams must be a positive integer, greater than 0')
  183. }
  184. if (useH2c != null && typeof useH2c !== 'boolean') {
  185. throw new InvalidArgumentError('useH2c must be a valid boolean value')
  186. }
  187. if (initialWindowSize != null && (!Number.isInteger(initialWindowSize) || initialWindowSize < 1)) {
  188. throw new InvalidArgumentError('initialWindowSize must be a positive integer, greater than 0')
  189. }
  190. if (connectionWindowSize != null && (!Number.isInteger(connectionWindowSize) || connectionWindowSize < 1)) {
  191. throw new InvalidArgumentError('connectionWindowSize must be a positive integer, greater than 0')
  192. }
  193. if (pingInterval != null && (typeof pingInterval !== 'number' || !Number.isInteger(pingInterval) || pingInterval < 0)) {
  194. throw new InvalidArgumentError('pingInterval must be a positive integer, greater or equal to 0')
  195. }
  196. super({ webSocket })
  197. if (typeof connect !== 'function') {
  198. connect = buildConnector({
  199. ...tls,
  200. maxCachedSessions,
  201. allowH2,
  202. useH2c,
  203. socketPath,
  204. timeout: connectTimeout,
  205. ...(typeof autoSelectFamily === 'boolean' ? { autoSelectFamily, autoSelectFamilyAttemptTimeout } : undefined),
  206. ...connect
  207. })
  208. } else {
  209. const customConnect = connect
  210. connect = (opts, callback) => customConnect({
  211. ...opts,
  212. ...(socketPath != null ? { socketPath } : null),
  213. ...(allowH2 != null ? { allowH2 } : null)
  214. }, callback)
  215. }
  216. this[kUrl] = util.parseOrigin(url)
  217. this[kConnector] = connect
  218. this[kPipelining] = pipelining != null ? pipelining : 1
  219. this[kMaxHeadersSize] = maxHeaderSize
  220. this[kKeepAliveDefaultTimeout] = keepAliveTimeout == null ? 4e3 : keepAliveTimeout
  221. this[kKeepAliveMaxTimeout] = keepAliveMaxTimeout == null ? 600e3 : keepAliveMaxTimeout
  222. this[kKeepAliveTimeoutThreshold] = keepAliveTimeoutThreshold == null ? 2e3 : keepAliveTimeoutThreshold
  223. this[kKeepAliveTimeoutValue] = this[kKeepAliveDefaultTimeout]
  224. this[kServerName] = null
  225. this[kLocalAddress] = localAddress != null ? localAddress : null
  226. this[kResuming] = 0 // 0, idle, 1, scheduled, 2 resuming
  227. this[kNeedDrain] = 0 // 0, idle, 1, scheduled, 2 resuming
  228. this[kHostHeader] = `host: ${this[kUrl].hostname}${this[kUrl].port ? `:${this[kUrl].port}` : ''}\r\n`
  229. this[kBodyTimeout] = bodyTimeout != null ? bodyTimeout : 300e3
  230. this[kHeadersTimeout] = headersTimeout != null ? headersTimeout : 300e3
  231. this[kStrictContentLength] = strictContentLength == null ? true : strictContentLength
  232. this[kMaxRequests] = maxRequestsPerClient
  233. this[kClosedResolve] = null
  234. this[kMaxResponseSize] = maxResponseSize > -1 ? maxResponseSize : -1
  235. this[kHTTPContext] = null
  236. // h2
  237. this[kMaxConcurrentStreams] = maxConcurrentStreams != null ? maxConcurrentStreams : 100 // Max peerConcurrentStreams for a Node h2 server
  238. // HTTP/2 window sizes are set to higher defaults than Node.js core for better performance:
  239. // - initialWindowSize: 262144 (256KB) vs Node.js default 65535 (64KB - 1)
  240. // Allows more data to be sent before requiring acknowledgment, improving throughput
  241. // especially on high-latency networks. This matches common production HTTP/2 servers.
  242. // - connectionWindowSize: 524288 (512KB) vs Node.js default (none set)
  243. // Provides better flow control for the entire connection across multiple streams.
  244. this[kHTTP2InitialWindowSize] = initialWindowSize != null ? initialWindowSize : 262144
  245. this[kHTTP2ConnectionWindowSize] = connectionWindowSize != null ? connectionWindowSize : 524288
  246. this[kPingInterval] = pingInterval != null ? pingInterval : 60e3 // Default ping interval for h2 - 1 minute
  247. // kQueue is built up of 3 sections separated by
  248. // the kRunningIdx and kPendingIdx indices.
  249. // | complete | running | pending |
  250. // ^ kRunningIdx ^ kPendingIdx ^ kQueue.length
  251. // kRunningIdx points to the first running element.
  252. // kPendingIdx points to the first pending element.
  253. // This implements a fast queue with an amortized
  254. // time of O(1).
  255. this[kQueue] = []
  256. this[kRunningIdx] = 0
  257. this[kPendingIdx] = 0
  258. this[kResume] = (sync) => resume(this, sync)
  259. this[kOnError] = (err) => onError(this, err)
  260. }
  261. get pipelining () {
  262. return this[kPipelining]
  263. }
  264. set pipelining (value) {
  265. this[kPipelining] = value
  266. this[kResume](true)
  267. }
  268. get stats () {
  269. return new ClientStats(this)
  270. }
  271. get [kPending] () {
  272. return this[kQueue].length - this[kPendingIdx]
  273. }
  274. get [kRunning] () {
  275. return this[kPendingIdx] - this[kRunningIdx]
  276. }
  277. get [kSize] () {
  278. return this[kQueue].length - this[kRunningIdx]
  279. }
  280. get [kConnected] () {
  281. return !!this[kHTTPContext] && !this[kConnecting] && !this[kHTTPContext].destroyed
  282. }
  283. get [kBusy] () {
  284. return Boolean(
  285. this[kHTTPContext]?.busy(null) ||
  286. (this[kSize] >= (getPipelining(this) || 1)) ||
  287. this[kPending] > 0
  288. )
  289. }
  290. [kConnect] (cb) {
  291. connect(this)
  292. this.once('connect', cb)
  293. }
  294. [kDispatch] (opts, handler) {
  295. const request = new Request(this[kUrl].origin, opts, handler)
  296. this[kQueue].push(request)
  297. if (this[kResuming]) {
  298. // Do nothing.
  299. } else if (util.bodyLength(request.body) == null && util.isIterable(request.body)) {
  300. // Wait a tick in case stream/iterator is ended in the same tick.
  301. this[kResuming] = 1
  302. queueMicrotask(() => resume(this))
  303. } else {
  304. this[kResume](true)
  305. }
  306. if (this[kResuming] && this[kNeedDrain] !== 2 && this[kBusy]) {
  307. this[kNeedDrain] = 2
  308. }
  309. return this[kNeedDrain] < 2
  310. }
  311. [kClose] () {
  312. // TODO: for H2 we need to gracefully flush the remaining enqueued
  313. // request and close each stream.
  314. return new Promise((resolve) => {
  315. if (this[kSize]) {
  316. this[kClosedResolve] = resolve
  317. } else {
  318. resolve(null)
  319. }
  320. })
  321. }
  322. [kDestroy] (err) {
  323. return new Promise((resolve) => {
  324. const requests = this[kQueue].splice(this[kPendingIdx])
  325. for (let i = 0; i < requests.length; i++) {
  326. const request = requests[i]
  327. if (request != null) {
  328. util.errorRequest(this, request, err)
  329. }
  330. }
  331. const callback = () => {
  332. if (this[kClosedResolve]) {
  333. // TODO (fix): Should we error here with ClientDestroyedError?
  334. this[kClosedResolve]()
  335. this[kClosedResolve] = null
  336. }
  337. resolve(null)
  338. }
  339. if (this[kHTTPContext]) {
  340. this[kHTTPContext].destroy(err, callback)
  341. this[kHTTPContext] = null
  342. } else {
  343. queueMicrotask(callback)
  344. }
  345. this[kResume]()
  346. })
  347. }
  348. }
  349. function onError (client, err) {
  350. if (
  351. client[kRunning] === 0 &&
  352. err.code !== 'UND_ERR_INFO' &&
  353. err.code !== 'UND_ERR_SOCKET'
  354. ) {
  355. // Error is not caused by running request and not a recoverable
  356. // socket error.
  357. assert(client[kPendingIdx] === client[kRunningIdx])
  358. const requests = client[kQueue].splice(client[kRunningIdx])
  359. for (let i = 0; i < requests.length; i++) {
  360. const request = requests[i]
  361. if (request != null) {
  362. util.errorRequest(client, request, err)
  363. }
  364. }
  365. assert(client[kSize] === 0)
  366. }
  367. }
  368. /**
  369. * @param {Client} client
  370. * @returns {void}
  371. */
  372. function connect (client) {
  373. assert(!client[kConnecting])
  374. assert(!client[kHTTPContext])
  375. let { host, hostname, protocol, port } = client[kUrl]
  376. // Resolve ipv6
  377. if (hostname[0] === '[') {
  378. const idx = hostname.indexOf(']')
  379. assert(idx !== -1)
  380. const ip = hostname.substring(1, idx)
  381. assert(net.isIPv6(ip))
  382. hostname = ip
  383. }
  384. client[kConnecting] = true
  385. if (channels.beforeConnect.hasSubscribers) {
  386. channels.beforeConnect.publish({
  387. connectParams: {
  388. host,
  389. hostname,
  390. protocol,
  391. port,
  392. version: client[kHTTPContext]?.version,
  393. servername: client[kServerName],
  394. localAddress: client[kLocalAddress]
  395. },
  396. connector: client[kConnector]
  397. })
  398. }
  399. try {
  400. client[kConnector]({
  401. host,
  402. hostname,
  403. protocol,
  404. port,
  405. servername: client[kServerName],
  406. localAddress: client[kLocalAddress]
  407. }, (err, socket) => {
  408. if (err) {
  409. handleConnectError(client, err, { host, hostname, protocol, port })
  410. client[kResume]()
  411. return
  412. }
  413. if (client.destroyed) {
  414. util.destroy(socket.on('error', noop), new ClientDestroyedError())
  415. client[kResume]()
  416. return
  417. }
  418. assert(socket)
  419. try {
  420. client[kHTTPContext] = socket.alpnProtocol === 'h2'
  421. ? connectH2(client, socket)
  422. : connectH1(client, socket)
  423. } catch (err) {
  424. socket.destroy().on('error', noop)
  425. handleConnectError(client, err, { host, hostname, protocol, port })
  426. client[kResume]()
  427. return
  428. }
  429. client[kConnecting] = false
  430. socket[kCounter] = 0
  431. socket[kMaxRequests] = client[kMaxRequests]
  432. socket[kClient] = client
  433. socket[kError] = null
  434. if (channels.connected.hasSubscribers) {
  435. channels.connected.publish({
  436. connectParams: {
  437. host,
  438. hostname,
  439. protocol,
  440. port,
  441. version: client[kHTTPContext]?.version,
  442. servername: client[kServerName],
  443. localAddress: client[kLocalAddress]
  444. },
  445. connector: client[kConnector],
  446. socket
  447. })
  448. }
  449. client.emit('connect', client[kUrl], [client])
  450. client[kResume]()
  451. })
  452. } catch (err) {
  453. handleConnectError(client, err, { host, hostname, protocol, port })
  454. client[kResume]()
  455. }
  456. }
  457. function handleConnectError (client, err, { host, hostname, protocol, port }) {
  458. if (client.destroyed) {
  459. return
  460. }
  461. client[kConnecting] = false
  462. if (channels.connectError.hasSubscribers) {
  463. channels.connectError.publish({
  464. connectParams: {
  465. host,
  466. hostname,
  467. protocol,
  468. port,
  469. version: client[kHTTPContext]?.version,
  470. servername: client[kServerName],
  471. localAddress: client[kLocalAddress]
  472. },
  473. connector: client[kConnector],
  474. error: err
  475. })
  476. }
  477. if (err.code === 'ERR_TLS_CERT_ALTNAME_INVALID') {
  478. assert(client[kRunning] === 0)
  479. while (client[kPending] > 0 && client[kQueue][client[kPendingIdx]].servername === client[kServerName]) {
  480. const request = client[kQueue][client[kPendingIdx]++]
  481. util.errorRequest(client, request, err)
  482. }
  483. } else {
  484. onError(client, err)
  485. }
  486. client.emit('connectionError', client[kUrl], [client], err)
  487. }
  488. function emitDrain (client) {
  489. client[kNeedDrain] = 0
  490. client.emit('drain', client[kUrl], [client])
  491. }
  492. function resume (client, sync) {
  493. if (client[kResuming] === 2) {
  494. return
  495. }
  496. client[kResuming] = 2
  497. _resume(client, sync)
  498. client[kResuming] = 0
  499. if (client[kRunningIdx] > 256) {
  500. client[kQueue].splice(0, client[kRunningIdx])
  501. client[kPendingIdx] -= client[kRunningIdx]
  502. client[kRunningIdx] = 0
  503. }
  504. }
  505. function _resume (client, sync) {
  506. while (true) {
  507. if (client.destroyed) {
  508. assert(client[kPending] === 0)
  509. return
  510. }
  511. if (client[kClosedResolve] && !client[kSize]) {
  512. client[kClosedResolve]()
  513. client[kClosedResolve] = null
  514. return
  515. }
  516. if (client[kHTTPContext]) {
  517. client[kHTTPContext].resume()
  518. }
  519. if (client[kBusy]) {
  520. client[kNeedDrain] = 2
  521. } else if (client[kNeedDrain] === 2) {
  522. if (sync) {
  523. client[kNeedDrain] = 1
  524. queueMicrotask(() => emitDrain(client))
  525. } else {
  526. emitDrain(client)
  527. }
  528. continue
  529. }
  530. if (client[kPending] === 0) {
  531. return
  532. }
  533. if (client[kRunning] >= (getPipelining(client) || 1)) {
  534. return
  535. }
  536. const request = client[kQueue][client[kPendingIdx]]
  537. if (request === null) {
  538. return
  539. }
  540. if (client[kUrl].protocol === 'https:' && client[kServerName] !== request.servername) {
  541. if (client[kRunning] > 0) {
  542. return
  543. }
  544. client[kServerName] = request.servername
  545. client[kHTTPContext]?.destroy(new InformationalError('servername changed'), () => {
  546. client[kHTTPContext] = null
  547. resume(client)
  548. })
  549. }
  550. if (client[kConnecting]) {
  551. return
  552. }
  553. if (!client[kHTTPContext]) {
  554. connect(client)
  555. return
  556. }
  557. if (client[kHTTPContext].destroyed) {
  558. return
  559. }
  560. if (client[kHTTPContext].busy(request)) {
  561. return
  562. }
  563. if (!request.aborted && client[kHTTPContext].write(request)) {
  564. client[kPendingIdx]++
  565. } else {
  566. client[kQueue].splice(client[kPendingIdx], 1)
  567. }
  568. }
  569. }
  570. module.exports = Client