client-h1.js 44 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746
  1. 'use strict'
  2. /* global WebAssembly */
  3. const assert = require('node:assert')
  4. const util = require('../core/util.js')
  5. const { channels } = require('../core/diagnostics.js')
  6. const timers = require('../util/timers.js')
  7. const {
  8. RequestContentLengthMismatchError,
  9. ResponseContentLengthMismatchError,
  10. RequestAbortedError,
  11. InvalidArgumentError,
  12. HeadersTimeoutError,
  13. HeadersOverflowError,
  14. SocketError,
  15. InformationalError,
  16. BodyTimeoutError,
  17. HTTPParserError,
  18. ResponseExceededMaxSizeError
  19. } = require('../core/errors.js')
  20. const {
  21. kUrl,
  22. kReset,
  23. kClient,
  24. kParser,
  25. kBlocking,
  26. kRunning,
  27. kPending,
  28. kSize,
  29. kWriting,
  30. kQueue,
  31. kNoRef,
  32. kKeepAliveDefaultTimeout,
  33. kHostHeader,
  34. kPendingIdx,
  35. kRunningIdx,
  36. kError,
  37. kPipelining,
  38. kSocket,
  39. kKeepAliveTimeoutValue,
  40. kMaxHeadersSize,
  41. kKeepAliveMaxTimeout,
  42. kKeepAliveTimeoutThreshold,
  43. kHeadersTimeout,
  44. kBodyTimeout,
  45. kStrictContentLength,
  46. kMaxRequests,
  47. kCounter,
  48. kMaxResponseSize,
  49. kOnError,
  50. kResume,
  51. kHTTPContext,
  52. kClosed
  53. } = require('../core/symbols.js')
  54. const constants = require('../llhttp/constants.js')
  55. const EMPTY_BUF = Buffer.alloc(0)
  56. const FastBuffer = Buffer[Symbol.species]
  57. const removeAllListeners = util.removeAllListeners
  58. const kIdleSocketValidation = Symbol('kIdleSocketValidation')
  59. const kIdleSocketValidationTimeout = Symbol('kIdleSocketValidationTimeout')
  60. const kSocketUsed = Symbol('kSocketUsed')
  61. let extractBody
  62. function lazyllhttp () {
  63. const llhttpWasmData = process.env.JEST_WORKER_ID ? require('../llhttp/llhttp-wasm.js') : undefined
  64. let mod
  65. // We disable wasm SIMD on older versions of Node.js on ppc64 that are broken on Power >=9 architectures.
  66. let useWasmSIMD = true
  67. if (process.arch === 'ppc64') {
  68. const [major, minor] = process.versions.node.split('.').map(n => parseInt(n, 10))
  69. if (major < 24 || (major === 24 && minor < 12)) {
  70. useWasmSIMD = false
  71. }
  72. }
  73. // The Env Variable UNDICI_NO_WASM_SIMD allows explicitly overriding the default behavior
  74. if (process.env.UNDICI_NO_WASM_SIMD === '1') {
  75. useWasmSIMD = false
  76. } else if (process.env.UNDICI_NO_WASM_SIMD === '0') {
  77. useWasmSIMD = true
  78. }
  79. if (useWasmSIMD) {
  80. try {
  81. mod = new WebAssembly.Module(require('../llhttp/llhttp_simd-wasm.js'))
  82. } catch {
  83. }
  84. }
  85. if (!mod) {
  86. // We could check if the error was caused by the simd option not
  87. // being enabled, but the occurring of this other error
  88. // * https://github.com/emscripten-core/emscripten/issues/11495
  89. // got me to remove that check to avoid breaking Node 12.
  90. mod = new WebAssembly.Module(llhttpWasmData || require('../llhttp/llhttp-wasm.js'))
  91. }
  92. return new WebAssembly.Instance(mod, {
  93. env: {
  94. /**
  95. * @param {number} p
  96. * @param {number} at
  97. * @param {number} len
  98. * @returns {number}
  99. */
  100. wasm_on_url: (p, at, len) => {
  101. return 0
  102. },
  103. /**
  104. * @param {number} p
  105. * @param {number} at
  106. * @param {number} len
  107. * @returns {number}
  108. */
  109. wasm_on_status: (p, at, len) => {
  110. assert(currentParser.ptr === p)
  111. const start = at - currentBufferPtr + currentBufferRef.byteOffset
  112. return currentParser.onStatus(new FastBuffer(currentBufferRef.buffer, start, len))
  113. },
  114. /**
  115. * @param {number} p
  116. * @returns {number}
  117. */
  118. wasm_on_message_begin: (p) => {
  119. assert(currentParser.ptr === p)
  120. return currentParser.onMessageBegin()
  121. },
  122. /**
  123. * @param {number} p
  124. * @param {number} at
  125. * @param {number} len
  126. * @returns {number}
  127. */
  128. wasm_on_header_field: (p, at, len) => {
  129. assert(currentParser.ptr === p)
  130. const start = at - currentBufferPtr + currentBufferRef.byteOffset
  131. return currentParser.onHeaderField(new FastBuffer(currentBufferRef.buffer, start, len))
  132. },
  133. /**
  134. * @param {number} p
  135. * @param {number} at
  136. * @param {number} len
  137. * @returns {number}
  138. */
  139. wasm_on_header_value: (p, at, len) => {
  140. assert(currentParser.ptr === p)
  141. const start = at - currentBufferPtr + currentBufferRef.byteOffset
  142. return currentParser.onHeaderValue(new FastBuffer(currentBufferRef.buffer, start, len))
  143. },
  144. /**
  145. * @param {number} p
  146. * @param {number} statusCode
  147. * @param {0|1} upgrade
  148. * @param {0|1} shouldKeepAlive
  149. * @returns {number}
  150. */
  151. wasm_on_headers_complete: (p, statusCode, upgrade, shouldKeepAlive) => {
  152. assert(currentParser.ptr === p)
  153. return currentParser.onHeadersComplete(statusCode, upgrade === 1, shouldKeepAlive === 1)
  154. },
  155. /**
  156. * @param {number} p
  157. * @param {number} at
  158. * @param {number} len
  159. * @returns {number}
  160. */
  161. wasm_on_body: (p, at, len) => {
  162. assert(currentParser.ptr === p)
  163. const start = at - currentBufferPtr + currentBufferRef.byteOffset
  164. return currentParser.onBody(new FastBuffer(currentBufferRef.buffer, start, len))
  165. },
  166. /**
  167. * @param {number} p
  168. * @returns {number}
  169. */
  170. wasm_on_message_complete: (p) => {
  171. assert(currentParser.ptr === p)
  172. return currentParser.onMessageComplete()
  173. }
  174. }
  175. })
  176. }
  177. let llhttpInstance = null
  178. /**
  179. * @type {Parser|null}
  180. */
  181. let currentParser = null
  182. let currentBufferRef = null
  183. /**
  184. * @type {number}
  185. */
  186. let currentBufferSize = 0
  187. let currentBufferPtr = null
  188. const USE_NATIVE_TIMER = 0
  189. const USE_FAST_TIMER = 1
  190. // Use fast timers for headers and body to take eventual event loop
  191. // latency into account.
  192. const TIMEOUT_HEADERS = 2 | USE_FAST_TIMER
  193. const TIMEOUT_BODY = 4 | USE_FAST_TIMER
  194. // Use native timers to ignore event loop latency for keep-alive
  195. // handling.
  196. const TIMEOUT_KEEP_ALIVE = 8 | USE_NATIVE_TIMER
  197. class Parser {
  198. /**
  199. * @param {import('./client.js')} client
  200. * @param {import('net').Socket} socket
  201. * @param {*} llhttp
  202. */
  203. constructor (client, socket, { exports }) {
  204. this.llhttp = exports
  205. this.ptr = this.llhttp.llhttp_alloc(constants.TYPE.RESPONSE)
  206. this.client = client
  207. /**
  208. * @type {import('net').Socket}
  209. */
  210. this.socket = socket
  211. this.timeout = null
  212. this.timeoutWeakRef = new WeakRef(this)
  213. this.timeoutValue = null
  214. this.timeoutType = null
  215. this.statusCode = 0
  216. this.statusText = ''
  217. this.upgrade = false
  218. this.headers = []
  219. this.headersSize = 0
  220. this.headersMaxSize = client[kMaxHeadersSize]
  221. this.shouldKeepAlive = false
  222. this.paused = false
  223. this.resume = this.resume.bind(this)
  224. this.bytesRead = 0
  225. this.keepAlive = ''
  226. this.contentLength = ''
  227. this.connection = ''
  228. this.maxResponseSize = client[kMaxResponseSize]
  229. }
  230. setTimeout (delay, type) {
  231. // If the existing timer and the new timer are of different timer type
  232. // (fast or native) or have different delay, we need to clear the existing
  233. // timer and set a new one.
  234. if (
  235. delay !== this.timeoutValue ||
  236. (type & USE_FAST_TIMER) ^ (this.timeoutType & USE_FAST_TIMER)
  237. ) {
  238. // If a timeout is already set, clear it with clearTimeout of the fast
  239. // timer implementation, as it can clear fast and native timers.
  240. if (this.timeout) {
  241. timers.clearTimeout(this.timeout)
  242. this.timeout = null
  243. }
  244. if (delay) {
  245. if (type & USE_FAST_TIMER) {
  246. this.timeout = timers.setFastTimeout(onParserTimeout, delay, this.timeoutWeakRef)
  247. } else {
  248. this.timeout = setTimeout(onParserTimeout, delay, this.timeoutWeakRef)
  249. this.timeout?.unref()
  250. }
  251. }
  252. this.timeoutValue = delay
  253. } else if (this.timeout) {
  254. if (this.timeout.refresh) {
  255. this.timeout.refresh()
  256. }
  257. }
  258. this.timeoutType = type
  259. }
  260. resume () {
  261. if (this.socket.destroyed || !this.paused) {
  262. return
  263. }
  264. assert(this.ptr != null)
  265. assert(currentParser === null)
  266. this.llhttp.llhttp_resume(this.ptr)
  267. assert(this.timeoutType === TIMEOUT_BODY)
  268. if (this.timeout) {
  269. if (this.timeout.refresh) {
  270. this.timeout.refresh()
  271. }
  272. }
  273. this.paused = false
  274. this.execute(this.socket.read() || EMPTY_BUF) // Flush parser.
  275. this.readMore()
  276. }
  277. readMore () {
  278. while (!this.paused && this.ptr) {
  279. const chunk = this.socket.read()
  280. if (chunk === null) {
  281. break
  282. }
  283. this.execute(chunk)
  284. }
  285. }
  286. /**
  287. * @param {Buffer} chunk
  288. */
  289. execute (chunk) {
  290. assert(currentParser === null)
  291. assert(this.ptr != null)
  292. assert(!this.paused)
  293. const { socket, llhttp } = this
  294. // Allocate a new buffer if the current buffer is too small.
  295. if (chunk.length > currentBufferSize) {
  296. if (currentBufferPtr) {
  297. llhttp.free(currentBufferPtr)
  298. }
  299. // Allocate a buffer that is a multiple of 4096 bytes.
  300. currentBufferSize = Math.ceil(chunk.length / 4096) * 4096
  301. currentBufferPtr = llhttp.malloc(currentBufferSize)
  302. }
  303. new Uint8Array(llhttp.memory.buffer, currentBufferPtr, currentBufferSize).set(chunk)
  304. // Call `execute` on the wasm parser.
  305. // We pass the `llhttp_parser` pointer address, the pointer address of buffer view data,
  306. // and finally the length of bytes to parse.
  307. // The return value is an error code or `constants.ERROR.OK`.
  308. try {
  309. let ret
  310. try {
  311. currentBufferRef = chunk
  312. currentParser = this
  313. ret = llhttp.llhttp_execute(this.ptr, currentBufferPtr, chunk.length)
  314. } finally {
  315. currentParser = null
  316. currentBufferRef = null
  317. }
  318. if (ret !== constants.ERROR.OK) {
  319. const data = chunk.subarray(llhttp.llhttp_get_error_pos(this.ptr) - currentBufferPtr)
  320. if (ret === constants.ERROR.PAUSED_UPGRADE) {
  321. this.onUpgrade(data)
  322. } else if (ret === constants.ERROR.PAUSED) {
  323. this.paused = true
  324. socket.unshift(data)
  325. } else {
  326. throw this.createError(ret, data)
  327. }
  328. }
  329. } catch (err) {
  330. util.destroy(socket, err)
  331. }
  332. }
  333. finish () {
  334. assert(currentParser === null)
  335. assert(this.ptr != null)
  336. assert(!this.paused)
  337. const { llhttp } = this
  338. let ret
  339. try {
  340. currentParser = this
  341. ret = llhttp.llhttp_finish(this.ptr)
  342. } finally {
  343. currentParser = null
  344. }
  345. if (ret === constants.ERROR.OK) {
  346. return null
  347. }
  348. if (ret === constants.ERROR.PAUSED || ret === constants.ERROR.PAUSED_UPGRADE) {
  349. this.paused = true
  350. return null
  351. }
  352. return this.createError(ret, EMPTY_BUF)
  353. }
  354. createError (ret, data) {
  355. const { llhttp, contentLength, bytesRead } = this
  356. if (contentLength && bytesRead !== parseInt(contentLength, 10)) {
  357. return new ResponseContentLengthMismatchError()
  358. }
  359. const ptr = llhttp.llhttp_get_error_reason(this.ptr)
  360. let message = ''
  361. if (ptr) {
  362. const len = new Uint8Array(llhttp.memory.buffer, ptr).indexOf(0)
  363. message =
  364. 'Response does not match the HTTP/1.1 protocol (' +
  365. Buffer.from(llhttp.memory.buffer, ptr, len).toString() +
  366. ')'
  367. }
  368. return new HTTPParserError(message, constants.ERROR[ret], data)
  369. }
  370. destroy () {
  371. assert(currentParser === null)
  372. assert(this.ptr != null)
  373. this.llhttp.llhttp_free(this.ptr)
  374. this.ptr = null
  375. this.timeout && timers.clearTimeout(this.timeout)
  376. this.timeout = null
  377. this.timeoutValue = null
  378. this.timeoutType = null
  379. this.paused = false
  380. }
  381. /**
  382. * @param {Buffer} buf
  383. * @returns {0}
  384. */
  385. onStatus (buf) {
  386. this.statusText = buf.toString()
  387. return 0
  388. }
  389. /**
  390. * @returns {0|-1}
  391. */
  392. onMessageBegin () {
  393. const { socket, client } = this
  394. if (socket.destroyed) {
  395. return -1
  396. }
  397. if (client[kRunning] === 0) {
  398. util.destroy(socket, new SocketError('bad response', util.getSocketInfo(socket)))
  399. return -1
  400. }
  401. const request = client[kQueue][client[kRunningIdx]]
  402. if (!request) {
  403. return -1
  404. }
  405. request.onResponseStarted()
  406. return 0
  407. }
  408. /**
  409. * @param {Buffer} buf
  410. * @returns {number}
  411. */
  412. onHeaderField (buf) {
  413. const len = this.headers.length
  414. if ((len & 1) === 0) {
  415. this.headers.push(buf)
  416. } else {
  417. this.headers[len - 1] = Buffer.concat([this.headers[len - 1], buf])
  418. }
  419. this.trackHeader(buf.length)
  420. return 0
  421. }
  422. /**
  423. * @param {Buffer} buf
  424. * @returns {number}
  425. */
  426. onHeaderValue (buf) {
  427. let len = this.headers.length
  428. if ((len & 1) === 1) {
  429. this.headers.push(buf)
  430. len += 1
  431. } else {
  432. this.headers[len - 1] = Buffer.concat([this.headers[len - 1], buf])
  433. }
  434. const key = this.headers[len - 2]
  435. if (key.length === 10) {
  436. const headerName = util.bufferToLowerCasedHeaderName(key)
  437. if (headerName === 'keep-alive') {
  438. this.keepAlive += buf.toString()
  439. } else if (headerName === 'connection') {
  440. this.connection += buf.toString()
  441. }
  442. } else if (key.length === 14 && util.bufferToLowerCasedHeaderName(key) === 'content-length') {
  443. this.contentLength += buf.toString()
  444. }
  445. this.trackHeader(buf.length)
  446. return 0
  447. }
  448. /**
  449. * @param {number} len
  450. */
  451. trackHeader (len) {
  452. this.headersSize += len
  453. if (this.headersSize >= this.headersMaxSize) {
  454. util.destroy(this.socket, new HeadersOverflowError())
  455. }
  456. }
  457. /**
  458. * @param {Buffer} head
  459. */
  460. onUpgrade (head) {
  461. const { upgrade, client, socket, headers, statusCode, statusText } = this
  462. assert(upgrade)
  463. assert(client[kSocket] === socket)
  464. assert(!socket.destroyed)
  465. assert(!this.paused)
  466. assert((headers.length & 1) === 0)
  467. const request = client[kQueue][client[kRunningIdx]]
  468. assert(request)
  469. assert(request.upgrade || request.method === 'CONNECT')
  470. this.statusCode = 0
  471. this.statusText = ''
  472. this.shouldKeepAlive = false
  473. this.headers = []
  474. this.headersSize = 0
  475. socket.unshift(head)
  476. socket[kParser].destroy()
  477. socket[kParser] = null
  478. socket[kClient] = null
  479. socket[kError] = null
  480. removeAllListeners(socket)
  481. client[kSocket] = null
  482. client[kHTTPContext] = null // TODO (fix): This is hacky...
  483. client[kQueue][client[kRunningIdx]++] = null
  484. client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade'))
  485. try {
  486. request.onUpgrade(statusCode, headers, socket, statusText)
  487. } catch (err) {
  488. util.errorRequest(client, request, err)
  489. util.destroy(socket, err)
  490. }
  491. client[kResume]()
  492. }
  493. /**
  494. * @param {number} statusCode
  495. * @param {boolean} upgrade
  496. * @param {boolean} shouldKeepAlive
  497. * @returns {number}
  498. */
  499. onHeadersComplete (statusCode, upgrade, shouldKeepAlive) {
  500. const { client, socket, headers, statusText } = this
  501. if (socket.destroyed) {
  502. return -1
  503. }
  504. if (client[kRunning] === 0) {
  505. util.destroy(socket, new SocketError('bad response', util.getSocketInfo(socket)))
  506. return -1
  507. }
  508. const request = client[kQueue][client[kRunningIdx]]
  509. if (!request) {
  510. return -1
  511. }
  512. assert(!this.upgrade)
  513. assert(this.statusCode < 200)
  514. if (statusCode === 100) {
  515. util.destroy(socket, new SocketError('bad response', util.getSocketInfo(socket)))
  516. return -1
  517. }
  518. /* this can only happen if server is misbehaving */
  519. if (upgrade && !request.upgrade) {
  520. util.destroy(socket, new SocketError('bad upgrade', util.getSocketInfo(socket)))
  521. return -1
  522. }
  523. assert(this.timeoutType === TIMEOUT_HEADERS)
  524. this.statusCode = statusCode
  525. this.shouldKeepAlive = (
  526. shouldKeepAlive ||
  527. // Override llhttp value which does not allow keepAlive for HEAD.
  528. (request.method === 'HEAD' && !socket[kReset] && this.connection.toLowerCase() === 'keep-alive')
  529. )
  530. if (this.statusCode >= 200) {
  531. const bodyTimeout = request.bodyTimeout != null
  532. ? request.bodyTimeout
  533. : client[kBodyTimeout]
  534. this.setTimeout(bodyTimeout, TIMEOUT_BODY)
  535. } else if (this.timeout) {
  536. if (this.timeout.refresh) {
  537. this.timeout.refresh()
  538. }
  539. }
  540. if (request.method === 'CONNECT') {
  541. assert(client[kRunning] === 1)
  542. this.upgrade = true
  543. return 2
  544. }
  545. if (upgrade) {
  546. assert(client[kRunning] === 1)
  547. this.upgrade = true
  548. return 2
  549. }
  550. assert((this.headers.length & 1) === 0)
  551. this.headers = []
  552. this.headersSize = 0
  553. if (this.shouldKeepAlive && client[kPipelining]) {
  554. const keepAliveTimeout = this.keepAlive ? util.parseKeepAliveTimeout(this.keepAlive) : null
  555. if (keepAliveTimeout != null) {
  556. const timeout = Math.min(
  557. keepAliveTimeout - client[kKeepAliveTimeoutThreshold],
  558. client[kKeepAliveMaxTimeout]
  559. )
  560. if (timeout <= 0) {
  561. socket[kReset] = true
  562. } else {
  563. client[kKeepAliveTimeoutValue] = timeout
  564. }
  565. } else {
  566. client[kKeepAliveTimeoutValue] = client[kKeepAliveDefaultTimeout]
  567. }
  568. } else {
  569. // Stop more requests from being dispatched.
  570. socket[kReset] = true
  571. }
  572. const pause = request.onHeaders(statusCode, headers, this.resume, statusText) === false
  573. if (request.aborted) {
  574. return -1
  575. }
  576. if (request.method === 'HEAD') {
  577. return 1
  578. }
  579. if (statusCode < 200) {
  580. return 1
  581. }
  582. if (socket[kBlocking]) {
  583. socket[kBlocking] = false
  584. client[kResume]()
  585. }
  586. return pause ? constants.ERROR.PAUSED : 0
  587. }
  588. /**
  589. * @param {Buffer} buf
  590. * @returns {number}
  591. */
  592. onBody (buf) {
  593. const { client, socket, statusCode, maxResponseSize } = this
  594. if (socket.destroyed) {
  595. return -1
  596. }
  597. const request = client[kQueue][client[kRunningIdx]]
  598. assert(request)
  599. assert(this.timeoutType === TIMEOUT_BODY)
  600. if (this.timeout) {
  601. if (this.timeout.refresh) {
  602. this.timeout.refresh()
  603. }
  604. }
  605. assert(statusCode >= 200)
  606. if (maxResponseSize > -1 && this.bytesRead + buf.length > maxResponseSize) {
  607. util.destroy(socket, new ResponseExceededMaxSizeError())
  608. return -1
  609. }
  610. this.bytesRead += buf.length
  611. if (request.onData(buf) === false) {
  612. return constants.ERROR.PAUSED
  613. }
  614. return 0
  615. }
  616. /**
  617. * @returns {number}
  618. */
  619. onMessageComplete () {
  620. const { client, socket, statusCode, upgrade, headers, contentLength, bytesRead, shouldKeepAlive } = this
  621. if (socket.destroyed && (!statusCode || shouldKeepAlive)) {
  622. return -1
  623. }
  624. if (upgrade) {
  625. return 0
  626. }
  627. assert(statusCode >= 100)
  628. assert((this.headers.length & 1) === 0)
  629. const request = client[kQueue][client[kRunningIdx]]
  630. assert(request)
  631. this.statusCode = 0
  632. this.statusText = ''
  633. this.bytesRead = 0
  634. this.contentLength = ''
  635. this.keepAlive = ''
  636. this.connection = ''
  637. this.headers = []
  638. this.headersSize = 0
  639. if (statusCode < 200) {
  640. return 0
  641. }
  642. if (request.method !== 'HEAD' && contentLength && bytesRead !== parseInt(contentLength, 10)) {
  643. util.destroy(socket, new ResponseContentLengthMismatchError())
  644. return -1
  645. }
  646. request.onComplete(headers)
  647. client[kQueue][client[kRunningIdx]++] = null
  648. socket[kSocketUsed] = client[kPending] === 0
  649. if (socket[kWriting]) {
  650. assert(client[kRunning] === 0)
  651. // Response completed before request.
  652. util.destroy(socket, new InformationalError('reset'))
  653. return constants.ERROR.PAUSED
  654. } else if (!shouldKeepAlive) {
  655. util.destroy(socket, new InformationalError('reset'))
  656. return constants.ERROR.PAUSED
  657. } else if (socket[kReset] && client[kRunning] === 0) {
  658. // Destroy socket once all requests have completed.
  659. // The request at the tail of the pipeline is the one
  660. // that requested reset and no further requests should
  661. // have been queued since then.
  662. util.destroy(socket, new InformationalError('reset'))
  663. return constants.ERROR.PAUSED
  664. } else if (client[kPipelining] == null || client[kPipelining] === 1) {
  665. // We must wait a full event loop cycle to reuse this socket to make sure
  666. // that non-spec compliant servers are not closing the connection even if they
  667. // said they won't.
  668. setImmediate(client[kResume])
  669. } else {
  670. client[kResume]()
  671. }
  672. return 0
  673. }
  674. }
  675. function onParserTimeout (parserWeakRef) {
  676. const parser = parserWeakRef.deref()
  677. if (!parser) {
  678. return
  679. }
  680. const { socket, timeoutType, client, paused } = parser
  681. if (timeoutType === TIMEOUT_HEADERS) {
  682. if (!socket[kWriting] || socket.writableNeedDrain || client[kRunning] > 1) {
  683. assert(!paused, 'cannot be paused while waiting for headers')
  684. util.destroy(socket, new HeadersTimeoutError())
  685. }
  686. } else if (timeoutType === TIMEOUT_BODY) {
  687. if (!paused) {
  688. util.destroy(socket, new BodyTimeoutError())
  689. }
  690. } else if (timeoutType === TIMEOUT_KEEP_ALIVE) {
  691. assert(client[kRunning] === 0 && client[kKeepAliveTimeoutValue])
  692. util.destroy(socket, new InformationalError('socket idle timeout'))
  693. }
  694. }
  695. /**
  696. * @param {import ('./client.js')} client
  697. * @param {import('net').Socket} socket
  698. * @returns
  699. */
  700. function connectH1 (client, socket) {
  701. client[kSocket] = socket
  702. if (!llhttpInstance) {
  703. llhttpInstance = lazyllhttp()
  704. }
  705. if (socket.errored) {
  706. throw socket.errored
  707. }
  708. if (socket.destroyed) {
  709. throw new SocketError('destroyed')
  710. }
  711. socket[kNoRef] = false
  712. socket[kWriting] = false
  713. socket[kReset] = false
  714. socket[kBlocking] = false
  715. socket[kIdleSocketValidation] = 0
  716. socket[kIdleSocketValidationTimeout] = null
  717. socket[kSocketUsed] = false
  718. socket[kParser] = new Parser(client, socket, llhttpInstance)
  719. util.addListener(socket, 'error', onHttpSocketError)
  720. util.addListener(socket, 'readable', onHttpSocketReadable)
  721. util.addListener(socket, 'end', onHttpSocketEnd)
  722. util.addListener(socket, 'close', onHttpSocketClose)
  723. socket[kClosed] = false
  724. socket.on('close', onSocketClose)
  725. return {
  726. version: 'h1',
  727. defaultPipelining: 1,
  728. write (request) {
  729. return writeH1(client, request)
  730. },
  731. resume () {
  732. resumeH1(client)
  733. },
  734. /**
  735. * @param {Error|undefined} err
  736. * @param {() => void} callback
  737. */
  738. destroy (err, callback) {
  739. if (socket[kClosed]) {
  740. queueMicrotask(callback)
  741. } else {
  742. socket.on('close', callback)
  743. socket.destroy(err)
  744. }
  745. },
  746. /**
  747. * @returns {boolean}
  748. */
  749. get destroyed () {
  750. return socket.destroyed
  751. },
  752. /**
  753. * @param {import('../core/request.js')} request
  754. * @returns {boolean}
  755. */
  756. busy (request) {
  757. if (socket[kWriting] || socket[kReset] || socket[kBlocking] || socket[kIdleSocketValidation] === 1) {
  758. return true
  759. }
  760. if (request) {
  761. if (client[kRunning] > 0 && !request.idempotent) {
  762. // Non-idempotent request cannot be retried.
  763. // Ensure that no other requests are inflight and
  764. // could cause failure.
  765. return true
  766. }
  767. if (client[kRunning] > 0 && (request.upgrade || request.method === 'CONNECT')) {
  768. // Don't dispatch an upgrade until all preceding requests have completed.
  769. // A misbehaving server might upgrade the connection before all pipelined
  770. // request has completed.
  771. return true
  772. }
  773. if (client[kRunning] > 0 && util.bodyLength(request.body) !== 0 &&
  774. (util.isStream(request.body) || util.isAsyncIterable(request.body) || util.isFormDataLike(request.body))) {
  775. // Request with stream or iterator body can error while other requests
  776. // are inflight and indirectly error those as well.
  777. // Ensure this doesn't happen by waiting for inflight
  778. // to complete before dispatching.
  779. // Request with stream or iterator body cannot be retried.
  780. // Ensure that no other requests are inflight and
  781. // could cause failure.
  782. return true
  783. }
  784. }
  785. return false
  786. }
  787. }
  788. }
  789. function onHttpSocketError (err) {
  790. assert(err.code !== 'ERR_TLS_CERT_ALTNAME_INVALID')
  791. const parser = this[kParser]
  792. // On Mac OS, we get an ECONNRESET even if there is a full body to be forwarded
  793. // to the user.
  794. if (err.code === 'ECONNRESET' && parser.statusCode && !parser.shouldKeepAlive) {
  795. const parserErr = parser.finish()
  796. if (parserErr) {
  797. this[kError] = parserErr
  798. this[kClient][kOnError](parserErr)
  799. }
  800. return
  801. }
  802. this[kError] = err
  803. this[kClient][kOnError](err)
  804. }
  805. function onHttpSocketReadable () {
  806. this[kParser]?.readMore()
  807. }
  808. function onHttpSocketEnd () {
  809. const parser = this[kParser]
  810. if (parser.statusCode && !parser.shouldKeepAlive) {
  811. const parserErr = parser.finish()
  812. if (parserErr) {
  813. util.destroy(this, parserErr)
  814. }
  815. return
  816. }
  817. util.destroy(this, new SocketError('other side closed', util.getSocketInfo(this)))
  818. }
  819. function onHttpSocketClose () {
  820. const parser = this[kParser]
  821. clearIdleSocketValidation(this)
  822. if (parser) {
  823. if (!this[kError] && parser.statusCode && !parser.shouldKeepAlive) {
  824. this[kError] = parser.finish() || this[kError]
  825. }
  826. this[kParser].destroy()
  827. this[kParser] = null
  828. }
  829. const err = this[kError] || new SocketError('closed', util.getSocketInfo(this))
  830. const client = this[kClient]
  831. client[kSocket] = null
  832. client[kHTTPContext] = null // TODO (fix): This is hacky...
  833. if (client.destroyed) {
  834. assert(client[kPending] === 0)
  835. // Fail entire queue.
  836. const requests = client[kQueue].splice(client[kRunningIdx])
  837. for (let i = 0; i < requests.length; i++) {
  838. const request = requests[i]
  839. util.errorRequest(client, request, err)
  840. }
  841. } else if (client[kRunning] > 0 && err.code !== 'UND_ERR_INFO') {
  842. // Fail head of pipeline.
  843. const request = client[kQueue][client[kRunningIdx]]
  844. client[kQueue][client[kRunningIdx]++] = null
  845. util.errorRequest(client, request, err)
  846. }
  847. client[kPendingIdx] = client[kRunningIdx]
  848. assert(client[kRunning] === 0)
  849. client.emit('disconnect', client[kUrl], [client], err)
  850. client[kResume]()
  851. }
  852. function onSocketClose () {
  853. this[kClosed] = true
  854. }
  855. function clearIdleSocketValidation (socket) {
  856. if (socket[kIdleSocketValidationTimeout]) {
  857. clearImmediate(socket[kIdleSocketValidationTimeout])
  858. socket[kIdleSocketValidationTimeout] = null
  859. }
  860. socket[kIdleSocketValidation] = 0
  861. }
  862. function scheduleIdleSocketValidation (client, socket) {
  863. socket[kIdleSocketValidation] = 1
  864. // Yield to the check phase (after poll) so unsolicited bytes / FIN / RST
  865. // already pending on this idle keep-alive socket are processed before the
  866. // next request is written (GHSA-35p6-xmwp-9g52).
  867. //
  868. // setTimeout(0) pays Node's ~1ms timer floor on every sequential reuse
  869. // (#5493). setImmediate avoids that, but an *unref'd* Immediate lets poll
  870. // block for ~500ms when the event loop is otherwise idle (#5600 / #5606).
  871. // A ref'd Immediate both keeps the pending request alive and makes poll
  872. // return immediately — the hybrid those issues asked for.
  873. socket[kIdleSocketValidationTimeout] = setImmediate(() => {
  874. socket[kIdleSocketValidationTimeout] = null
  875. socket[kIdleSocketValidation] = 2
  876. if (client[kSocket] === socket && !socket.destroyed) {
  877. client[kResume]()
  878. }
  879. })
  880. }
  881. /**
  882. * @param {import('./client.js')} client
  883. */
  884. function resumeH1 (client) {
  885. const socket = client[kSocket]
  886. if (socket && !socket.destroyed) {
  887. if (client[kSize] === 0) {
  888. if (!socket[kNoRef] && socket.unref) {
  889. socket.unref()
  890. socket[kNoRef] = true
  891. }
  892. } else if (socket[kNoRef] && socket.ref) {
  893. socket.ref()
  894. socket[kNoRef] = false
  895. }
  896. if (client[kRunning] === 0 && client[kPending] > 0 && socket[kSocketUsed]) {
  897. if (socket[kIdleSocketValidation] === 0) {
  898. scheduleIdleSocketValidation(client, socket)
  899. socket[kParser].readMore()
  900. if (socket.destroyed) {
  901. return
  902. }
  903. return
  904. }
  905. if (socket[kIdleSocketValidation] === 1) {
  906. socket[kParser].readMore()
  907. if (socket.destroyed) {
  908. return
  909. }
  910. return
  911. }
  912. }
  913. if (client[kRunning] === 0) {
  914. socket[kParser].readMore()
  915. if (socket.destroyed) {
  916. return
  917. }
  918. }
  919. if (client[kSize] === 0) {
  920. if (socket[kParser].timeoutType !== TIMEOUT_KEEP_ALIVE) {
  921. socket[kParser].setTimeout(client[kKeepAliveTimeoutValue], TIMEOUT_KEEP_ALIVE)
  922. }
  923. } else if (client[kRunning] > 0 && socket[kParser].statusCode < 200) {
  924. if (socket[kParser].timeoutType !== TIMEOUT_HEADERS) {
  925. const request = client[kQueue][client[kRunningIdx]]
  926. const headersTimeout = request.headersTimeout != null
  927. ? request.headersTimeout
  928. : client[kHeadersTimeout]
  929. socket[kParser].setTimeout(headersTimeout, TIMEOUT_HEADERS)
  930. }
  931. }
  932. }
  933. }
  934. // https://www.rfc-editor.org/rfc/rfc7230#section-3.3.2
  935. function shouldSendContentLength (method) {
  936. return method !== 'GET' && method !== 'HEAD' && method !== 'OPTIONS' && method !== 'TRACE' && method !== 'CONNECT'
  937. }
  938. /**
  939. * @param {import('./client.js')} client
  940. * @param {import('../core/request.js')} request
  941. * @returns
  942. */
  943. function writeH1 (client, request) {
  944. const { method, path, host, upgrade, blocking, reset } = request
  945. let { body, headers, contentLength } = request
  946. // https://tools.ietf.org/html/rfc7231#section-4.3.1
  947. // https://tools.ietf.org/html/rfc7231#section-4.3.2
  948. // https://tools.ietf.org/html/rfc7231#section-4.3.5
  949. // Sending a payload body on a request that does not
  950. // expect it can cause undefined behavior on some
  951. // servers and corrupt connection state. Do not
  952. // re-use the connection for further requests.
  953. const expectsPayload = (
  954. method === 'PUT' ||
  955. method === 'POST' ||
  956. method === 'PATCH' ||
  957. method === 'QUERY' ||
  958. method === 'PROPFIND' ||
  959. method === 'PROPPATCH'
  960. )
  961. if (util.isFormDataLike(body)) {
  962. if (!extractBody) {
  963. extractBody = require('../web/fetch/body.js').extractBody
  964. }
  965. const [bodyStream, contentType] = extractBody(body)
  966. if (request.contentType == null) {
  967. headers.push('content-type', contentType)
  968. }
  969. body = bodyStream.stream
  970. contentLength = bodyStream.length
  971. } else if (util.isBlobLike(body) && request.contentType == null) {
  972. const contentType = body.type
  973. if (contentType) {
  974. const contentTypeValue = `${contentType}`
  975. if (!util.isValidHeaderValue(contentTypeValue)) {
  976. util.errorRequest(client, request, new InvalidArgumentError('invalid content-type header'))
  977. return false
  978. }
  979. headers.push('content-type', contentTypeValue)
  980. }
  981. }
  982. if (body && typeof body.read === 'function') {
  983. // Try to read EOF in order to get length.
  984. body.read(0)
  985. }
  986. const bodyLength = util.bodyLength(body)
  987. contentLength = bodyLength ?? contentLength
  988. if (contentLength === null) {
  989. contentLength = request.contentLength
  990. }
  991. if (contentLength === 0 && !expectsPayload) {
  992. // https://tools.ietf.org/html/rfc7230#section-3.3.2
  993. // A user agent SHOULD NOT send a Content-Length header field when
  994. // the request message does not contain a payload body and the method
  995. // semantics do not anticipate such a body.
  996. contentLength = null
  997. }
  998. // https://github.com/nodejs/undici/issues/2046
  999. // A user agent may send a Content-Length header with 0 value, this should be allowed.
  1000. if (shouldSendContentLength(method) && contentLength > 0 && request.contentLength !== null && request.contentLength !== contentLength) {
  1001. if (client[kStrictContentLength]) {
  1002. util.errorRequest(client, request, new RequestContentLengthMismatchError())
  1003. return false
  1004. }
  1005. process.emitWarning(new RequestContentLengthMismatchError())
  1006. }
  1007. const socket = client[kSocket]
  1008. clearIdleSocketValidation(socket)
  1009. /**
  1010. * @param {Error} [err]
  1011. * @returns {void}
  1012. */
  1013. const abort = (err) => {
  1014. if (request.aborted || request.completed) {
  1015. return
  1016. }
  1017. util.errorRequest(client, request, err || new RequestAbortedError())
  1018. util.destroy(body)
  1019. util.destroy(socket, new InformationalError('aborted'))
  1020. }
  1021. try {
  1022. request.onConnect(abort)
  1023. } catch (err) {
  1024. util.errorRequest(client, request, err)
  1025. }
  1026. if (request.aborted) {
  1027. return false
  1028. }
  1029. if (method === 'HEAD') {
  1030. // https://github.com/mcollina/undici/issues/258
  1031. // Close after a HEAD request to interop with misbehaving servers
  1032. // that may send a body in the response.
  1033. socket[kReset] = true
  1034. }
  1035. if (upgrade || method === 'CONNECT') {
  1036. // On CONNECT or upgrade, block pipeline from dispatching further
  1037. // requests on this connection.
  1038. socket[kReset] = true
  1039. }
  1040. if (reset != null) {
  1041. socket[kReset] = reset
  1042. }
  1043. if (client[kMaxRequests] && socket[kCounter]++ >= client[kMaxRequests]) {
  1044. socket[kReset] = true
  1045. }
  1046. if (blocking) {
  1047. socket[kBlocking] = true
  1048. }
  1049. if (socket.setTypeOfService) {
  1050. socket.setTypeOfService(request.typeOfService)
  1051. }
  1052. let header = `${method} ${path} HTTP/1.1\r\n`
  1053. if (typeof host === 'string') {
  1054. header += `host: ${host}\r\n`
  1055. } else {
  1056. header += client[kHostHeader]
  1057. }
  1058. if (upgrade) {
  1059. header += `connection: upgrade\r\nupgrade: ${upgrade}\r\n`
  1060. } else if (client[kPipelining] && !socket[kReset]) {
  1061. header += 'connection: keep-alive\r\n'
  1062. } else {
  1063. header += 'connection: close\r\n'
  1064. }
  1065. if (Array.isArray(headers)) {
  1066. for (let n = 0; n < headers.length; n += 2) {
  1067. const key = headers[n + 0]
  1068. const val = headers[n + 1]
  1069. if (Array.isArray(val)) {
  1070. for (let i = 0; i < val.length; i++) {
  1071. header += `${key}: ${val[i]}\r\n`
  1072. }
  1073. } else {
  1074. header += `${key}: ${val}\r\n`
  1075. }
  1076. }
  1077. }
  1078. if (channels.sendHeaders.hasSubscribers) {
  1079. channels.sendHeaders.publish({ request, headers: header, socket })
  1080. }
  1081. if (!body || bodyLength === 0) {
  1082. writeBuffer(abort, null, client, request, socket, contentLength, header, expectsPayload)
  1083. } else if (util.isBuffer(body)) {
  1084. writeBuffer(abort, body, client, request, socket, contentLength, header, expectsPayload)
  1085. } else if (util.isBlobLike(body)) {
  1086. if (typeof body.stream === 'function') {
  1087. writeIterable(abort, body.stream(), client, request, socket, contentLength, header, expectsPayload)
  1088. } else {
  1089. writeBlob(abort, body, client, request, socket, contentLength, header, expectsPayload)
  1090. }
  1091. } else if (util.isStream(body)) {
  1092. writeStream(abort, body, client, request, socket, contentLength, header, expectsPayload)
  1093. } else if (util.isIterable(body)) {
  1094. writeIterable(abort, body, client, request, socket, contentLength, header, expectsPayload)
  1095. } else {
  1096. assert(false)
  1097. }
  1098. return true
  1099. }
  1100. /**
  1101. * @param {AbortCallback} abort
  1102. * @param {import('stream').Stream} body
  1103. * @param {import('./client.js')} client
  1104. * @param {import('../core/request.js')} request
  1105. * @param {import('net').Socket} socket
  1106. * @param {number} contentLength
  1107. * @param {string} header
  1108. * @param {boolean} expectsPayload
  1109. */
  1110. function writeStream (abort, body, client, request, socket, contentLength, header, expectsPayload) {
  1111. assert(contentLength !== 0 || client[kRunning] === 0, 'stream body cannot be pipelined')
  1112. let finished = false
  1113. const writer = new AsyncWriter({ abort, socket, request, contentLength, client, expectsPayload, header })
  1114. /**
  1115. * @param {Buffer} chunk
  1116. * @returns {void}
  1117. */
  1118. const onData = function (chunk) {
  1119. if (finished) {
  1120. return
  1121. }
  1122. try {
  1123. if (!writer.write(chunk) && this.pause) {
  1124. this.pause()
  1125. }
  1126. } catch (err) {
  1127. util.destroy(this, err)
  1128. }
  1129. }
  1130. /**
  1131. * @returns {void}
  1132. */
  1133. const onDrain = function () {
  1134. if (finished) {
  1135. return
  1136. }
  1137. if (body.resume) {
  1138. body.resume()
  1139. }
  1140. }
  1141. /**
  1142. * @returns {void}
  1143. */
  1144. const onClose = function () {
  1145. // 'close' might be emitted *before* 'error' for
  1146. // broken streams. Wait a tick to avoid this case.
  1147. queueMicrotask(() => {
  1148. // It's only safe to remove 'error' listener after
  1149. // 'close'.
  1150. body.removeListener('error', onFinished)
  1151. })
  1152. if (!finished) {
  1153. const err = new RequestAbortedError()
  1154. queueMicrotask(() => onFinished(err))
  1155. }
  1156. }
  1157. /**
  1158. * @param {Error} [err]
  1159. * @returns
  1160. */
  1161. const onFinished = function (err) {
  1162. if (finished) {
  1163. return
  1164. }
  1165. finished = true
  1166. assert(socket.destroyed || (socket[kWriting] && client[kRunning] <= 1))
  1167. socket
  1168. .off('drain', onDrain)
  1169. .off('error', onFinished)
  1170. body
  1171. .removeListener('data', onData)
  1172. .removeListener('end', onFinished)
  1173. .removeListener('close', onClose)
  1174. if (!err) {
  1175. try {
  1176. writer.end()
  1177. } catch (er) {
  1178. err = er
  1179. }
  1180. }
  1181. writer.destroy(err)
  1182. if (err && (err.code !== 'UND_ERR_INFO' || err.message !== 'reset')) {
  1183. util.destroy(body, err)
  1184. } else {
  1185. util.destroy(body)
  1186. }
  1187. }
  1188. body
  1189. .on('data', onData)
  1190. .on('end', onFinished)
  1191. .on('error', onFinished)
  1192. .on('close', onClose)
  1193. if (body.resume) {
  1194. body.resume()
  1195. }
  1196. socket
  1197. .on('drain', onDrain)
  1198. .on('error', onFinished)
  1199. if (body.errorEmitted ?? body.errored) {
  1200. setImmediate(onFinished, body.errored)
  1201. } else if (body.endEmitted ?? body.readableEnded) {
  1202. setImmediate(onFinished, null)
  1203. }
  1204. if (body.closeEmitted ?? body.closed) {
  1205. setImmediate(onClose)
  1206. }
  1207. }
  1208. /**
  1209. * @typedef AbortCallback
  1210. * @type {Function}
  1211. * @param {Error} [err]
  1212. * @returns {void}
  1213. */
  1214. /**
  1215. * @param {AbortCallback} abort
  1216. * @param {Uint8Array|null} body
  1217. * @param {import('./client.js')} client
  1218. * @param {import('../core/request.js')} request
  1219. * @param {import('net').Socket} socket
  1220. * @param {number} contentLength
  1221. * @param {string} header
  1222. * @param {boolean} expectsPayload
  1223. * @returns {void}
  1224. */
  1225. function writeBuffer (abort, body, client, request, socket, contentLength, header, expectsPayload) {
  1226. try {
  1227. if (!body) {
  1228. if (contentLength === 0) {
  1229. socket.write(`${header}content-length: 0\r\n\r\n`, 'latin1')
  1230. } else {
  1231. assert(contentLength === null, 'no body must not have content length')
  1232. socket.write(`${header}\r\n`, 'latin1')
  1233. }
  1234. } else if (util.isBuffer(body)) {
  1235. assert(contentLength === body.byteLength, 'buffer body must have content length')
  1236. socket.cork()
  1237. socket.write(`${header}content-length: ${contentLength}\r\n\r\n`, 'latin1')
  1238. socket.write(body)
  1239. socket.uncork()
  1240. request.onBodySent(body)
  1241. if (!expectsPayload && request.reset !== false) {
  1242. socket[kReset] = true
  1243. }
  1244. }
  1245. request.onRequestSent()
  1246. client[kResume]()
  1247. } catch (err) {
  1248. abort(err)
  1249. }
  1250. }
  1251. /**
  1252. * @param {AbortCallback} abort
  1253. * @param {Blob} body
  1254. * @param {import('./client.js')} client
  1255. * @param {import('../core/request.js')} request
  1256. * @param {import('net').Socket} socket
  1257. * @param {number} contentLength
  1258. * @param {string} header
  1259. * @param {boolean} expectsPayload
  1260. * @returns {Promise<void>}
  1261. */
  1262. async function writeBlob (abort, body, client, request, socket, contentLength, header, expectsPayload) {
  1263. assert(contentLength === body.size, 'blob body must have content length')
  1264. try {
  1265. if (contentLength != null && contentLength !== body.size) {
  1266. throw new RequestContentLengthMismatchError()
  1267. }
  1268. const buffer = Buffer.from(await body.arrayBuffer())
  1269. socket.cork()
  1270. socket.write(`${header}content-length: ${contentLength}\r\n\r\n`, 'latin1')
  1271. socket.write(buffer)
  1272. socket.uncork()
  1273. request.onBodySent(buffer)
  1274. request.onRequestSent()
  1275. if (!expectsPayload && request.reset !== false) {
  1276. socket[kReset] = true
  1277. }
  1278. client[kResume]()
  1279. } catch (err) {
  1280. abort(err)
  1281. }
  1282. }
  1283. /**
  1284. * @param {AbortCallback} abort
  1285. * @param {Iterable} body
  1286. * @param {import('./client.js')} client
  1287. * @param {import('../core/request.js')} request
  1288. * @param {import('net').Socket} socket
  1289. * @param {number} contentLength
  1290. * @param {string} header
  1291. * @param {boolean} expectsPayload
  1292. * @returns {Promise<void>}
  1293. */
  1294. async function writeIterable (abort, body, client, request, socket, contentLength, header, expectsPayload) {
  1295. assert(contentLength !== 0 || client[kRunning] === 0, 'iterator body cannot be pipelined')
  1296. let callback = null
  1297. function onDrain () {
  1298. if (callback) {
  1299. const cb = callback
  1300. callback = null
  1301. cb()
  1302. }
  1303. }
  1304. const waitForDrain = () => new Promise((resolve, reject) => {
  1305. assert(callback === null)
  1306. if (socket[kError]) {
  1307. reject(socket[kError])
  1308. } else {
  1309. callback = resolve
  1310. }
  1311. })
  1312. socket
  1313. .on('close', onDrain)
  1314. .on('drain', onDrain)
  1315. const writer = new AsyncWriter({ abort, socket, request, contentLength, client, expectsPayload, header })
  1316. try {
  1317. // It's up to the user to somehow abort the async iterable.
  1318. for await (const chunk of body) {
  1319. if (socket[kError]) {
  1320. throw socket[kError]
  1321. }
  1322. if (!writer.write(chunk)) {
  1323. await waitForDrain()
  1324. }
  1325. }
  1326. writer.end()
  1327. } catch (err) {
  1328. writer.destroy(err)
  1329. } finally {
  1330. socket
  1331. .off('close', onDrain)
  1332. .off('drain', onDrain)
  1333. }
  1334. }
  1335. class AsyncWriter {
  1336. /**
  1337. *
  1338. * @param {object} arg
  1339. * @param {AbortCallback} arg.abort
  1340. * @param {import('net').Socket} arg.socket
  1341. * @param {import('../core/request.js')} arg.request
  1342. * @param {number} arg.contentLength
  1343. * @param {import('./client.js')} arg.client
  1344. * @param {boolean} arg.expectsPayload
  1345. * @param {string} arg.header
  1346. */
  1347. constructor ({ abort, socket, request, contentLength, client, expectsPayload, header }) {
  1348. this.socket = socket
  1349. this.request = request
  1350. this.contentLength = contentLength
  1351. this.client = client
  1352. this.bytesWritten = 0
  1353. this.expectsPayload = expectsPayload
  1354. this.header = header
  1355. this.abort = abort
  1356. socket[kWriting] = true
  1357. }
  1358. /**
  1359. * @param {Buffer} chunk
  1360. * @returns
  1361. */
  1362. write (chunk) {
  1363. const { socket, request, contentLength, client, bytesWritten, expectsPayload, header } = this
  1364. if (socket[kError]) {
  1365. throw socket[kError]
  1366. }
  1367. if (socket.destroyed) {
  1368. return false
  1369. }
  1370. const len = Buffer.byteLength(chunk)
  1371. if (!len) {
  1372. return true
  1373. }
  1374. // We should defer writing chunks.
  1375. if (contentLength !== null && bytesWritten + len > contentLength) {
  1376. if (client[kStrictContentLength]) {
  1377. throw new RequestContentLengthMismatchError()
  1378. }
  1379. process.emitWarning(new RequestContentLengthMismatchError())
  1380. }
  1381. socket.cork()
  1382. if (bytesWritten === 0) {
  1383. if (!expectsPayload && request.reset !== false) {
  1384. socket[kReset] = true
  1385. }
  1386. if (contentLength === null) {
  1387. socket.write(`${header}transfer-encoding: chunked\r\n`, 'latin1')
  1388. } else {
  1389. socket.write(`${header}content-length: ${contentLength}\r\n\r\n`, 'latin1')
  1390. }
  1391. }
  1392. if (contentLength === null) {
  1393. socket.write(`\r\n${len.toString(16)}\r\n`, 'latin1')
  1394. }
  1395. this.bytesWritten += len
  1396. const ret = socket.write(chunk)
  1397. socket.uncork()
  1398. request.onBodySent(chunk)
  1399. if (!ret) {
  1400. if (socket[kParser].timeout && socket[kParser].timeoutType === TIMEOUT_HEADERS) {
  1401. if (socket[kParser].timeout.refresh) {
  1402. socket[kParser].timeout.refresh()
  1403. }
  1404. }
  1405. }
  1406. return ret
  1407. }
  1408. /**
  1409. * @returns {void}
  1410. */
  1411. end () {
  1412. const { socket, contentLength, client, bytesWritten, expectsPayload, header, request } = this
  1413. request.onRequestSent()
  1414. socket[kWriting] = false
  1415. if (socket[kError]) {
  1416. throw socket[kError]
  1417. }
  1418. if (socket.destroyed) {
  1419. return
  1420. }
  1421. if (bytesWritten === 0) {
  1422. if (expectsPayload) {
  1423. // https://tools.ietf.org/html/rfc7230#section-3.3.2
  1424. // A user agent SHOULD send a Content-Length in a request message when
  1425. // no Transfer-Encoding is sent and the request method defines a meaning
  1426. // for an enclosed payload body.
  1427. socket.write(`${header}content-length: 0\r\n\r\n`, 'latin1')
  1428. } else {
  1429. socket.write(`${header}\r\n`, 'latin1')
  1430. }
  1431. } else if (contentLength === null) {
  1432. socket.write('\r\n0\r\n\r\n', 'latin1')
  1433. }
  1434. if (contentLength !== null && bytesWritten !== contentLength) {
  1435. if (client[kStrictContentLength]) {
  1436. throw new RequestContentLengthMismatchError()
  1437. } else {
  1438. process.emitWarning(new RequestContentLengthMismatchError())
  1439. }
  1440. }
  1441. if (socket[kParser].timeout && socket[kParser].timeoutType === TIMEOUT_HEADERS) {
  1442. if (socket[kParser].timeout.refresh) {
  1443. socket[kParser].timeout.refresh()
  1444. }
  1445. }
  1446. client[kResume]()
  1447. }
  1448. /**
  1449. * @param {Error} [err]
  1450. * @returns {void}
  1451. */
  1452. destroy (err) {
  1453. const { socket, client, abort } = this
  1454. socket[kWriting] = false
  1455. if (err) {
  1456. assert(client[kRunning] <= 1, 'pipeline should only contain this request')
  1457. abort(err)
  1458. }
  1459. }
  1460. }
  1461. module.exports = connectH1