client-h2.js 29 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060
  1. 'use strict'
  2. const assert = require('node:assert')
  3. const { pipeline } = require('node:stream')
  4. const util = require('../core/util.js')
  5. const {
  6. RequestContentLengthMismatchError,
  7. RequestAbortedError,
  8. SocketError,
  9. InformationalError,
  10. InvalidArgumentError,
  11. HeadersTimeoutError,
  12. BodyTimeoutError
  13. } = require('../core/errors.js')
  14. const {
  15. kUrl,
  16. kReset,
  17. kClient,
  18. kRunning,
  19. kPending,
  20. kQueue,
  21. kPendingIdx,
  22. kRunningIdx,
  23. kError,
  24. kSocket,
  25. kStrictContentLength,
  26. kOnError,
  27. kMaxConcurrentStreams,
  28. kPingInterval,
  29. kHTTP2Session,
  30. kHTTP2InitialWindowSize,
  31. kHTTP2ConnectionWindowSize,
  32. kResume,
  33. kSize,
  34. kHTTPContext,
  35. kClosed,
  36. kBodyTimeout,
  37. kHeadersTimeout,
  38. kEnableConnectProtocol,
  39. kRemoteSettings,
  40. kHTTP2Stream,
  41. kHTTP2SessionState
  42. } = require('../core/symbols.js')
  43. const { channels } = require('../core/diagnostics.js')
  44. const kOpenStreams = Symbol('open streams')
  45. let extractBody
  46. /** @type {import('http2')} */
  47. let http2
  48. try {
  49. http2 = require('node:http2')
  50. } catch {
  51. // @ts-ignore
  52. http2 = { constants: {} }
  53. }
  54. const {
  55. constants: {
  56. HTTP2_HEADER_AUTHORITY,
  57. HTTP2_HEADER_METHOD,
  58. HTTP2_HEADER_PATH,
  59. HTTP2_HEADER_SCHEME,
  60. HTTP2_HEADER_CONTENT_LENGTH,
  61. HTTP2_HEADER_EXPECT,
  62. HTTP2_HEADER_STATUS,
  63. HTTP2_HEADER_PROTOCOL,
  64. NGHTTP2_REFUSED_STREAM,
  65. NGHTTP2_CANCEL
  66. }
  67. } = http2
  68. function parseH2Headers (headers) {
  69. const result = []
  70. for (const [name, value] of Object.entries(headers)) {
  71. // h2 may concat the header value by array
  72. // e.g. Set-Cookie
  73. if (Array.isArray(value)) {
  74. for (const subvalue of value) {
  75. // we need to provide each header value of header name
  76. // because the headers handler expect name-value pair
  77. result.push(Buffer.from(name), Buffer.from(subvalue))
  78. }
  79. } else {
  80. result.push(Buffer.from(name), Buffer.from(value))
  81. }
  82. }
  83. return result
  84. }
  85. function connectH2 (client, socket) {
  86. client[kSocket] = socket
  87. const http2InitialWindowSize = client[kHTTP2InitialWindowSize]
  88. const http2ConnectionWindowSize = client[kHTTP2ConnectionWindowSize]
  89. const session = http2.connect(client[kUrl], {
  90. createConnection: () => socket,
  91. peerMaxConcurrentStreams: client[kMaxConcurrentStreams],
  92. settings: {
  93. // TODO(metcoder95): add support for PUSH
  94. enablePush: false,
  95. ...(http2InitialWindowSize != null ? { initialWindowSize: http2InitialWindowSize } : null)
  96. }
  97. })
  98. client[kSocket] = socket
  99. session[kOpenStreams] = 0
  100. session[kClient] = client
  101. session[kSocket] = socket
  102. session[kHTTP2SessionState] = {
  103. ping: {
  104. interval: client[kPingInterval] === 0 ? null : setInterval(onHttp2SendPing, client[kPingInterval], session).unref()
  105. }
  106. }
  107. // We set it to true by default in a best-effort; however once connected to an H2 server
  108. // we will check if extended CONNECT protocol is supported or not
  109. // and set this value accordingly.
  110. session[kEnableConnectProtocol] = false
  111. // States whether or not we have received the remote settings from the server
  112. session[kRemoteSettings] = false
  113. // Apply connection-level flow control once connected (if supported).
  114. if (http2ConnectionWindowSize) {
  115. util.addListener(session, 'connect', applyConnectionWindowSize.bind(session, http2ConnectionWindowSize))
  116. }
  117. util.addListener(session, 'error', onHttp2SessionError)
  118. util.addListener(session, 'frameError', onHttp2FrameError)
  119. util.addListener(session, 'end', onHttp2SessionEnd)
  120. util.addListener(session, 'goaway', onHttp2SessionGoAway)
  121. util.addListener(session, 'close', onHttp2SessionClose)
  122. util.addListener(session, 'remoteSettings', onHttp2RemoteSettings)
  123. // TODO (@metcoder95): implement SETTINGS support
  124. // util.addListener(session, 'localSettings', onHttp2RemoteSettings)
  125. session.unref()
  126. client[kHTTP2Session] = session
  127. socket[kHTTP2Session] = session
  128. util.addListener(socket, 'error', onHttp2SocketError)
  129. util.addListener(socket, 'end', onHttp2SocketEnd)
  130. util.addListener(socket, 'close', onHttp2SocketClose)
  131. socket[kClosed] = false
  132. socket.on('close', onSocketClose)
  133. return {
  134. version: 'h2',
  135. defaultPipelining: Infinity,
  136. /**
  137. * @param {import('../core/request.js')} request
  138. * @returns {boolean}
  139. */
  140. write (request) {
  141. return writeH2(client, request)
  142. },
  143. /**
  144. * @returns {void}
  145. */
  146. resume () {
  147. resumeH2(client)
  148. },
  149. /**
  150. * @param {Error | null} err
  151. * @param {() => void} callback
  152. */
  153. destroy (err, callback) {
  154. if (socket[kClosed]) {
  155. queueMicrotask(callback)
  156. } else {
  157. socket.destroy(err).on('close', callback)
  158. }
  159. },
  160. /**
  161. * @type {boolean}
  162. */
  163. get destroyed () {
  164. return socket.destroyed
  165. },
  166. /**
  167. * @param {import('../core/request.js')} request
  168. * @returns {boolean}
  169. */
  170. busy (request) {
  171. if (request != null) {
  172. if (client[kRunning] > 0) {
  173. // We are already processing requests
  174. // Non-idempotent request cannot be retried.
  175. // Ensure that no other requests are inflight and
  176. // could cause failure.
  177. if (request.idempotent === false) return true
  178. // Don't dispatch an upgrade until all preceding requests have completed.
  179. // Possibly, we do not have remote settings confirmed yet.
  180. if ((request.upgrade === 'websocket' || request.method === 'CONNECT') && session[kRemoteSettings] === false) return true
  181. // Request with stream or iterator body can error while other requests
  182. // are inflight and indirectly error those as well.
  183. // Ensure this doesn't happen by waiting for inflight
  184. // to complete before dispatching.
  185. // Request with stream or iterator body cannot be retried.
  186. // Ensure that no other requests are inflight and
  187. // could cause failure.
  188. if (util.bodyLength(request.body) !== 0 &&
  189. (util.isStream(request.body) || util.isAsyncIterable(request.body) || util.isFormDataLike(request.body))) return true
  190. } else {
  191. return (request.upgrade === 'websocket' || request.method === 'CONNECT') && session[kRemoteSettings] === false
  192. }
  193. }
  194. return false
  195. }
  196. }
  197. }
  198. function resumeH2 (client) {
  199. const socket = client[kSocket]
  200. if (socket?.destroyed === false) {
  201. // Only let the process exit when there is genuinely nothing outstanding.
  202. // Unreffing because the peer advertised MAX_CONCURRENT_STREAMS = 0 left
  203. // queued requests with nothing holding the event loop open, so the process
  204. // could exit with status 0 while an awaited request never settled.
  205. if (client[kSize] === 0) {
  206. socket.unref()
  207. client[kHTTP2Session].unref()
  208. } else {
  209. socket.ref()
  210. client[kHTTP2Session].ref()
  211. }
  212. }
  213. }
  214. function applyConnectionWindowSize (connectionWindowSize) {
  215. try {
  216. if (typeof this.setLocalWindowSize === 'function') {
  217. this.setLocalWindowSize(connectionWindowSize)
  218. }
  219. } catch {
  220. // Best-effort only.
  221. }
  222. }
  223. function onHttp2RemoteSettings (settings) {
  224. // Fallbacks are a safe bet, remote setting will always override
  225. this[kClient][kMaxConcurrentStreams] = settings.maxConcurrentStreams ?? this[kClient][kMaxConcurrentStreams]
  226. /**
  227. * From RFC-8441
  228. * A sender MUST NOT send a SETTINGS_ENABLE_CONNECT_PROTOCOL parameter
  229. * with the value of 0 after previously sending a value of 1.
  230. */
  231. // Note: Cannot be tested in Node, it does not supports disabling the extended CONNECT protocol once enabled
  232. if (this[kRemoteSettings] === true && this[kEnableConnectProtocol] === true && settings.enableConnectProtocol === false) {
  233. const err = new InformationalError('HTTP/2: Server disabled extended CONNECT protocol against RFC-8441')
  234. this[kSocket][kError] = err
  235. this[kClient][kOnError](err)
  236. return
  237. }
  238. this[kEnableConnectProtocol] = settings.enableConnectProtocol ?? this[kEnableConnectProtocol]
  239. this[kRemoteSettings] = true
  240. this[kClient][kResume]()
  241. }
  242. function onHttp2SendPing (session) {
  243. const state = session[kHTTP2SessionState]
  244. if ((session.closed || session.destroyed) && state.ping.interval != null) {
  245. clearInterval(state.ping.interval)
  246. state.ping.interval = null
  247. return
  248. }
  249. // If no ping sent, do nothing
  250. session.ping(onPing.bind(session))
  251. function onPing (err, duration) {
  252. const client = this[kClient]
  253. const socket = this[kClient]
  254. if (err != null) {
  255. const error = new InformationalError(`HTTP/2: "PING" errored - type ${err.message}`)
  256. socket[kError] = error
  257. client[kOnError](error)
  258. } else {
  259. client.emit('ping', duration)
  260. }
  261. }
  262. }
  263. function onHttp2SessionError (err) {
  264. assert(err.code !== 'ERR_TLS_CERT_ALTNAME_INVALID')
  265. this[kSocket][kError] = err
  266. this[kClient][kOnError](err)
  267. }
  268. function onHttp2FrameError (type, code, id) {
  269. if (id === 0) {
  270. const err = new InformationalError(`HTTP/2: "frameError" received - type ${type}, code ${code}`)
  271. this[kSocket][kError] = err
  272. this[kClient][kOnError](err)
  273. }
  274. }
  275. function onHttp2SessionEnd () {
  276. const err = new SocketError('other side closed', util.getSocketInfo(this[kSocket]))
  277. this.destroy(err)
  278. util.destroy(this[kSocket], err)
  279. }
  280. /**
  281. * This is the root cause of #3011
  282. * We need to handle GOAWAY frames properly, and trigger the session close
  283. * along with the socket right away
  284. *
  285. * @this {import('http2').ClientHttp2Session}
  286. * @param {number} errorCode
  287. */
  288. // Backport of #5410 and #5569. HTTP/2 multiplexes, so requests complete out of
  289. // order; advancing kRunningIdx blindly retired whichever request happened to
  290. // sit at the head instead of the one that actually finished, which both lost
  291. // requests and left phantom running slots behind.
  292. function completeRequest (client, request, resetPendingIdx = false) {
  293. const queue = client[kQueue]
  294. const runningIdx = client[kRunningIdx]
  295. // In-order completion: clear the request and advance without splicing.
  296. // The client's resume loop compacts cleared slots once the index grows.
  297. if (runningIdx < client[kPendingIdx] && queue[runningIdx] === request) {
  298. queue[runningIdx] = null
  299. client[kRunningIdx] = runningIdx + 1
  300. return
  301. }
  302. const index = queue.indexOf(request, runningIdx)
  303. if (index === -1 || index >= client[kPendingIdx]) {
  304. return
  305. }
  306. queue.splice(index, 1)
  307. client[kPendingIdx]--
  308. if (resetPendingIdx && client[kPendingIdx] < client[kRunningIdx]) {
  309. client[kPendingIdx] = client[kRunningIdx]
  310. }
  311. }
  312. function onHttp2SessionGoAway (errorCode) {
  313. // TODO(mcollina): Verify if GOAWAY implements the spec correctly:
  314. // https://datatracker.ietf.org/doc/html/rfc7540#section-6.8
  315. // Specifically, we do not verify the "valid" stream id.
  316. const err = this[kError] || new SocketError(`HTTP/2: "GOAWAY" frame received with code ${errorCode}`, util.getSocketInfo(this[kSocket]))
  317. const client = this[kClient]
  318. client[kSocket] = null
  319. client[kHTTPContext] = null
  320. // this is an HTTP2 session
  321. this.close()
  322. this[kHTTP2Session] = null
  323. util.destroy(this[kSocket], err)
  324. // Fail head of pipeline.
  325. if (client[kRunningIdx] < client[kQueue].length) {
  326. const request = client[kQueue][client[kRunningIdx]]
  327. client[kQueue][client[kRunningIdx]++] = null
  328. if (request != null) {
  329. util.errorRequest(client, request, err)
  330. }
  331. client[kPendingIdx] = client[kRunningIdx]
  332. }
  333. assert(client[kRunning] === 0)
  334. client.emit('disconnect', client[kUrl], [client], err)
  335. client.emit('connectionError', client[kUrl], [client], err)
  336. client[kResume]()
  337. }
  338. function onHttp2SessionClose () {
  339. const { [kClient]: client, [kHTTP2SessionState]: state } = this
  340. const { [kSocket]: socket } = client
  341. const err = this[kSocket][kError] || this[kError] || new SocketError('closed', util.getSocketInfo(socket))
  342. client[kSocket] = null
  343. client[kHTTPContext] = null
  344. if (state.ping.interval != null) {
  345. clearInterval(state.ping.interval)
  346. state.ping.interval = null
  347. }
  348. if (client.destroyed) {
  349. assert(client[kPending] === 0)
  350. // Fail entire queue.
  351. const requests = client[kQueue].splice(client[kRunningIdx])
  352. for (let i = 0; i < requests.length; i++) {
  353. const request = requests[i]
  354. if (request != null) {
  355. util.errorRequest(client, request, err)
  356. }
  357. }
  358. }
  359. }
  360. function onHttp2SocketClose () {
  361. const err = this[kError] || new SocketError('closed', util.getSocketInfo(this))
  362. const client = this[kHTTP2Session][kClient]
  363. client[kSocket] = null
  364. client[kHTTPContext] = null
  365. if (this[kHTTP2Session] !== null) {
  366. this[kHTTP2Session].destroy(err)
  367. }
  368. client[kPendingIdx] = client[kRunningIdx]
  369. assert(client[kRunning] === 0)
  370. client.emit('disconnect', client[kUrl], [client], err)
  371. client[kResume]()
  372. }
  373. function onHttp2SocketError (err) {
  374. assert(err.code !== 'ERR_TLS_CERT_ALTNAME_INVALID')
  375. this[kError] = err
  376. this[kClient][kOnError](err)
  377. }
  378. function onHttp2SocketEnd () {
  379. util.destroy(this, new SocketError('other side closed', util.getSocketInfo(this)))
  380. }
  381. function onSocketClose () {
  382. this[kClosed] = true
  383. }
  384. // https://www.rfc-editor.org/rfc/rfc7230#section-3.3.2
  385. function shouldSendContentLength (method) {
  386. return method !== 'GET' && method !== 'HEAD' && method !== 'OPTIONS' && method !== 'TRACE' && method !== 'CONNECT'
  387. }
  388. function writeH2 (client, request) {
  389. // Time to the response headers, then time between body chunks. Using
  390. // bodyTimeout for both made headersTimeout a no-op over HTTP/2.
  391. const headersTimeout = request.headersTimeout ?? client[kHeadersTimeout]
  392. const bodyTimeout = request.bodyTimeout ?? client[kBodyTimeout]
  393. const session = client[kHTTP2Session]
  394. const { method, path, host, upgrade, expectContinue, signal, protocol, headers: reqHeaders } = request
  395. let { body } = request
  396. if (upgrade != null && upgrade !== 'websocket') {
  397. util.errorRequest(client, request, new InvalidArgumentError(`Custom upgrade "${upgrade}" not supported over HTTP/2`))
  398. return false
  399. }
  400. const headers = {}
  401. for (let n = 0; n < reqHeaders.length; n += 2) {
  402. const key = reqHeaders[n + 0]
  403. const val = reqHeaders[n + 1]
  404. if (key === 'cookie') {
  405. if (headers[key] != null) {
  406. headers[key] = Array.isArray(headers[key]) ? (headers[key].push(val), headers[key]) : [headers[key], val]
  407. } else {
  408. headers[key] = val
  409. }
  410. continue
  411. }
  412. if (Array.isArray(val)) {
  413. for (let i = 0; i < val.length; i++) {
  414. if (headers[key]) {
  415. headers[key] += `, ${val[i]}`
  416. } else {
  417. headers[key] = val[i]
  418. }
  419. }
  420. } else if (headers[key]) {
  421. headers[key] += `, ${val}`
  422. } else {
  423. headers[key] = val
  424. }
  425. }
  426. /** @type {import('node:http2').ClientHttp2Stream} */
  427. let stream = null
  428. const { hostname, port } = client[kUrl]
  429. headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}`
  430. headers[HTTP2_HEADER_METHOD] = method
  431. const abort = (err) => {
  432. if (request.aborted || request.completed) {
  433. return
  434. }
  435. err = err || new RequestAbortedError()
  436. util.errorRequest(client, request, err)
  437. if (stream != null) {
  438. // Some chunks might still come after abort,
  439. // let's ignore them
  440. stream.removeAllListeners('data')
  441. // On Abort, we close the stream to send RST_STREAM frame
  442. stream.close()
  443. // We move the running index to the next request
  444. client[kOnError](err)
  445. completeRequest(client, request)
  446. client[kResume]()
  447. }
  448. // We do not destroy the socket as we can continue using the session
  449. // the stream gets destroyed and the session remains to create new streams
  450. util.destroy(body, err)
  451. }
  452. try {
  453. // We are already connected, streams are pending.
  454. // We can call on connect, and wait for abort
  455. request.onConnect(abort)
  456. } catch (err) {
  457. util.errorRequest(client, request, err)
  458. }
  459. if (request.aborted) {
  460. return false
  461. }
  462. if (upgrade || method === 'CONNECT') {
  463. session.ref()
  464. if (upgrade === 'websocket') {
  465. // We cannot upgrade to websocket if extended CONNECT protocol is not supported
  466. if (session[kEnableConnectProtocol] === false) {
  467. util.errorRequest(client, request, new InformationalError('HTTP/2: Extended CONNECT protocol not supported by server'))
  468. session.unref()
  469. return false
  470. }
  471. // We force the method to CONNECT
  472. // as per RFC-8441
  473. // https://datatracker.ietf.org/doc/html/rfc8441#section-4
  474. headers[HTTP2_HEADER_METHOD] = 'CONNECT'
  475. headers[HTTP2_HEADER_PROTOCOL] = 'websocket'
  476. // :path and :scheme headers must be omitted when sending CONNECT but set if extended-CONNECT
  477. headers[HTTP2_HEADER_PATH] = path
  478. if (protocol === 'ws:' || protocol === 'wss:') {
  479. headers[HTTP2_HEADER_SCHEME] = protocol === 'ws:' ? 'http' : 'https'
  480. } else {
  481. headers[HTTP2_HEADER_SCHEME] = protocol === 'http:' ? 'http' : 'https'
  482. }
  483. stream = session.request(headers, { endStream: false, signal })
  484. stream[kHTTP2Stream] = true
  485. ++session[kOpenStreams]
  486. stream.once('response', (headers, _flags) => {
  487. const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers
  488. try {
  489. request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream)
  490. } catch (err) {
  491. abort(err)
  492. return
  493. }
  494. completeRequest(client, request)
  495. })
  496. stream.on('error', () => {
  497. if (stream.rstCode === NGHTTP2_REFUSED_STREAM || stream.rstCode === NGHTTP2_CANCEL) {
  498. // NGHTTP2_REFUSED_STREAM (7) or NGHTTP2_CANCEL (8)
  499. // We do not treat those as errors as the server might
  500. // not support websockets and refuse the stream
  501. abort(new InformationalError(`HTTP/2: "stream error" received - code ${stream.rstCode}`))
  502. }
  503. })
  504. stream.once('close', () => {
  505. session[kOpenStreams] -= 1
  506. if (session[kOpenStreams] === 0) session.unref()
  507. })
  508. stream.setTimeout(headersTimeout)
  509. return true
  510. }
  511. // TODO: consolidate once we support CONNECT properly
  512. // NOTE: We are already connected, streams are pending, first request
  513. // will create a new stream. We trigger a request to create the stream and wait until
  514. // `ready` event is triggered
  515. // We disabled endStream to allow the user to write to the stream
  516. stream = session.request(headers, { endStream: false, signal })
  517. stream[kHTTP2Stream] = true
  518. ++session[kOpenStreams]
  519. stream.on('response', headers => {
  520. const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers
  521. try {
  522. request.onUpgrade(statusCode, parseH2Headers(realHeaders), stream)
  523. } catch (err) {
  524. abort(err)
  525. return
  526. }
  527. completeRequest(client, request)
  528. })
  529. stream.on('error', abort)
  530. stream.once('close', () => {
  531. session[kOpenStreams] -= 1
  532. if (session[kOpenStreams] === 0) session.unref()
  533. })
  534. stream.setTimeout(headersTimeout)
  535. return true
  536. }
  537. // https://tools.ietf.org/html/rfc7540#section-8.3
  538. // :path and :scheme headers must be omitted when sending CONNECT
  539. headers[HTTP2_HEADER_PATH] = path
  540. headers[HTTP2_HEADER_SCHEME] = protocol === 'http:' ? 'http' : 'https'
  541. // https://tools.ietf.org/html/rfc7231#section-4.3.1
  542. // https://tools.ietf.org/html/rfc7231#section-4.3.2
  543. // https://tools.ietf.org/html/rfc7231#section-4.3.5
  544. // Sending a payload body on a request that does not
  545. // expect it can cause undefined behavior on some
  546. // servers and corrupt connection state. Do not
  547. // re-use the connection for further requests.
  548. const expectsPayload = (
  549. method === 'PUT' ||
  550. method === 'POST' ||
  551. method === 'PATCH'
  552. )
  553. if (body && typeof body.read === 'function') {
  554. // Try to read EOF in order to get length.
  555. body.read(0)
  556. }
  557. let contentLength = util.bodyLength(body)
  558. if (util.isFormDataLike(body)) {
  559. extractBody ??= require('../web/fetch/body.js').extractBody
  560. const [bodyStream, contentType] = extractBody(body)
  561. headers['content-type'] = contentType
  562. body = bodyStream.stream
  563. contentLength = bodyStream.length
  564. }
  565. if (contentLength == null) {
  566. contentLength = request.contentLength
  567. }
  568. if (!expectsPayload) {
  569. // https://tools.ietf.org/html/rfc7230#section-3.3.2
  570. // A user agent SHOULD NOT send a Content-Length header field when
  571. // the request message does not contain a payload body and the method
  572. // semantics do not anticipate such a body.
  573. // And for methods that don't expect a payload, omit Content-Length.
  574. contentLength = null
  575. }
  576. // https://github.com/nodejs/undici/issues/2046
  577. // A user agent may send a Content-Length header with 0 value, this should be allowed.
  578. if (shouldSendContentLength(method) && contentLength > 0 && request.contentLength != null && request.contentLength !== contentLength) {
  579. if (client[kStrictContentLength]) {
  580. util.errorRequest(client, request, new RequestContentLengthMismatchError())
  581. return false
  582. }
  583. process.emitWarning(new RequestContentLengthMismatchError())
  584. }
  585. if (contentLength != null) {
  586. assert(body || contentLength === 0, 'no body must not have content length')
  587. headers[HTTP2_HEADER_CONTENT_LENGTH] = `${contentLength}`
  588. }
  589. session.ref()
  590. if (channels.sendHeaders.hasSubscribers) {
  591. let header = ''
  592. for (const key in headers) {
  593. header += `${key}: ${headers[key]}\r\n`
  594. }
  595. channels.sendHeaders.publish({ request, headers: header, socket: session[kSocket] })
  596. }
  597. // TODO(metcoder95): add support for sending trailers
  598. const shouldEndStream = method === 'GET' || method === 'HEAD' || body === null
  599. if (expectContinue) {
  600. headers[HTTP2_HEADER_EXPECT] = '100-continue'
  601. stream = session.request(headers, { endStream: shouldEndStream, signal })
  602. stream[kHTTP2Stream] = true
  603. stream.once('continue', writeBodyH2)
  604. } else {
  605. stream = session.request(headers, {
  606. endStream: shouldEndStream,
  607. signal
  608. })
  609. stream[kHTTP2Stream] = true
  610. writeBodyH2()
  611. }
  612. // Increment counter as we have new streams open
  613. ++session[kOpenStreams]
  614. stream.setTimeout(headersTimeout)
  615. // Track whether we received a response (headers)
  616. let responseReceived = false
  617. stream.once('response', headers => {
  618. const { [HTTP2_HEADER_STATUS]: statusCode, ...realHeaders } = headers
  619. request.onResponseStarted()
  620. responseReceived = true
  621. stream.setTimeout(bodyTimeout)
  622. // Due to the stream nature, it is possible we face a race condition
  623. // where the stream has been assigned, but the request has been aborted
  624. // the request remains in-flight and headers hasn't been received yet
  625. // for those scenarios, best effort is to destroy the stream immediately
  626. // as there's no value to keep it open.
  627. if (request.aborted) {
  628. stream.removeAllListeners('data')
  629. return
  630. }
  631. if (request.onHeaders(Number(statusCode), parseH2Headers(realHeaders), stream.resume.bind(stream), '') === false) {
  632. stream.pause()
  633. }
  634. stream.on('data', (chunk) => {
  635. if (request.aborted || request.completed) {
  636. return
  637. }
  638. if (request.onData(chunk) === false) {
  639. stream.pause()
  640. }
  641. })
  642. })
  643. stream.once('end', () => {
  644. stream.removeAllListeners('data')
  645. // If we received a response, this is a normal completion
  646. if (responseReceived) {
  647. if (!request.aborted && !request.completed) {
  648. request.onComplete({})
  649. }
  650. completeRequest(client, request)
  651. client[kResume]()
  652. } else {
  653. // Stream ended without receiving a response - this is an error
  654. // (e.g., server destroyed the stream before sending headers)
  655. abort(new InformationalError('HTTP/2: stream half-closed (remote)'))
  656. completeRequest(client, request, true)
  657. client[kResume]()
  658. }
  659. })
  660. stream.once('close', () => {
  661. stream.removeAllListeners('data')
  662. session[kOpenStreams] -= 1
  663. if (session[kOpenStreams] === 0) {
  664. session.unref()
  665. }
  666. // A stream can close without ever emitting 'end' or 'error': a peer's
  667. // RST_STREAM(CANCEL) received before the response is reported by Node as a
  668. // bare 'close', and destroying the stream unenrolls its timeout, so no
  669. // 'timeout' follows either. Nothing else would ever settle this request.
  670. if (!request.aborted && !request.completed) {
  671. abort(new InformationalError('HTTP/2: stream closed before the response was complete'))
  672. }
  673. })
  674. stream.once('error', function (err) {
  675. stream.removeAllListeners('data')
  676. abort(err)
  677. })
  678. stream.once('frameError', (type, code) => {
  679. stream.removeAllListeners('data')
  680. abort(new InformationalError(`HTTP/2: "frameError" received - type ${type}, code ${code}`))
  681. })
  682. stream.on('aborted', () => {
  683. stream.removeAllListeners('data')
  684. })
  685. stream.on('timeout', () => {
  686. const err = responseReceived
  687. ? new BodyTimeoutError(`HTTP/2: "body timeout after ${bodyTimeout}"`)
  688. : new HeadersTimeoutError(`HTTP/2: "headers timeout after ${headersTimeout}"`)
  689. stream.removeAllListeners('data')
  690. session[kOpenStreams] -= 1
  691. if (session[kOpenStreams] === 0) {
  692. session.unref()
  693. }
  694. abort(err)
  695. })
  696. stream.once('trailers', trailers => {
  697. if (request.aborted || request.completed) {
  698. return
  699. }
  700. stream.removeAllListeners('data')
  701. request.onComplete(trailers)
  702. })
  703. return true
  704. function writeBodyH2 () {
  705. if (!body || contentLength === 0) {
  706. writeBuffer(
  707. abort,
  708. stream,
  709. null,
  710. client,
  711. request,
  712. client[kSocket],
  713. contentLength,
  714. expectsPayload
  715. )
  716. } else if (util.isBuffer(body)) {
  717. writeBuffer(
  718. abort,
  719. stream,
  720. body,
  721. client,
  722. request,
  723. client[kSocket],
  724. contentLength,
  725. expectsPayload
  726. )
  727. } else if (util.isBlobLike(body)) {
  728. if (typeof body.stream === 'function') {
  729. writeIterable(
  730. abort,
  731. stream,
  732. body.stream(),
  733. client,
  734. request,
  735. client[kSocket],
  736. contentLength,
  737. expectsPayload
  738. )
  739. } else {
  740. writeBlob(
  741. abort,
  742. stream,
  743. body,
  744. client,
  745. request,
  746. client[kSocket],
  747. contentLength,
  748. expectsPayload
  749. )
  750. }
  751. } else if (util.isStream(body)) {
  752. writeStream(
  753. abort,
  754. client[kSocket],
  755. expectsPayload,
  756. stream,
  757. body,
  758. client,
  759. request,
  760. contentLength
  761. )
  762. } else if (util.isIterable(body)) {
  763. writeIterable(
  764. abort,
  765. stream,
  766. body,
  767. client,
  768. request,
  769. client[kSocket],
  770. contentLength,
  771. expectsPayload
  772. )
  773. } else {
  774. assert(false)
  775. }
  776. }
  777. }
  778. function writeBuffer (abort, h2stream, body, client, request, socket, contentLength, expectsPayload) {
  779. try {
  780. if (body != null && util.isBuffer(body)) {
  781. assert(contentLength === body.byteLength, 'buffer body must have content length')
  782. h2stream.cork()
  783. h2stream.write(body)
  784. h2stream.uncork()
  785. h2stream.end()
  786. request.onBodySent(body)
  787. }
  788. if (!expectsPayload) {
  789. socket[kReset] = true
  790. }
  791. request.onRequestSent()
  792. client[kResume]()
  793. } catch (error) {
  794. abort(error)
  795. }
  796. }
  797. function writeStream (abort, socket, expectsPayload, h2stream, body, client, request, contentLength) {
  798. assert(contentLength !== 0 || client[kRunning] === 0, 'stream body cannot be pipelined')
  799. // For HTTP/2, is enough to pipe the stream
  800. const pipe = pipeline(
  801. body,
  802. h2stream,
  803. (err) => {
  804. if (err) {
  805. util.destroy(pipe, err)
  806. abort(err)
  807. } else {
  808. util.removeAllListeners(pipe)
  809. request.onRequestSent()
  810. if (!expectsPayload) {
  811. socket[kReset] = true
  812. }
  813. client[kResume]()
  814. }
  815. }
  816. )
  817. util.addListener(pipe, 'data', onPipeData)
  818. function onPipeData (chunk) {
  819. request.onBodySent(chunk)
  820. }
  821. }
  822. async function writeBlob (abort, h2stream, body, client, request, socket, contentLength, expectsPayload) {
  823. assert(contentLength === body.size, 'blob body must have content length')
  824. try {
  825. if (contentLength != null && contentLength !== body.size) {
  826. throw new RequestContentLengthMismatchError()
  827. }
  828. const buffer = Buffer.from(await body.arrayBuffer())
  829. h2stream.cork()
  830. h2stream.write(buffer)
  831. h2stream.uncork()
  832. h2stream.end()
  833. request.onBodySent(buffer)
  834. request.onRequestSent()
  835. if (!expectsPayload) {
  836. socket[kReset] = true
  837. }
  838. client[kResume]()
  839. } catch (err) {
  840. abort(err)
  841. }
  842. }
  843. async function writeIterable (abort, h2stream, body, client, request, socket, contentLength, expectsPayload) {
  844. assert(contentLength !== 0 || client[kRunning] === 0, 'iterator body cannot be pipelined')
  845. let callback = null
  846. function onDrain () {
  847. if (callback) {
  848. const cb = callback
  849. callback = null
  850. cb()
  851. }
  852. }
  853. const waitForDrain = () => new Promise((resolve, reject) => {
  854. assert(callback === null)
  855. if (socket[kError]) {
  856. reject(socket[kError])
  857. } else {
  858. callback = resolve
  859. }
  860. })
  861. h2stream
  862. .on('close', onDrain)
  863. .on('drain', onDrain)
  864. try {
  865. // It's up to the user to somehow abort the async iterable.
  866. for await (const chunk of body) {
  867. if (socket[kError]) {
  868. throw socket[kError]
  869. }
  870. const res = h2stream.write(chunk)
  871. request.onBodySent(chunk)
  872. if (!res) {
  873. await waitForDrain()
  874. }
  875. }
  876. h2stream.end()
  877. request.onRequestSent()
  878. if (!expectsPayload) {
  879. socket[kReset] = true
  880. }
  881. client[kResume]()
  882. } catch (err) {
  883. abort(err)
  884. } finally {
  885. h2stream
  886. .off('close', onDrain)
  887. .off('drain', onDrain)
  888. }
  889. }
  890. module.exports = connectH2