| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746 |
- 'use strict'
- /* global WebAssembly */
- const assert = require('node:assert')
- const util = require('../core/util.js')
- const { channels } = require('../core/diagnostics.js')
- const timers = require('../util/timers.js')
- const {
- RequestContentLengthMismatchError,
- ResponseContentLengthMismatchError,
- RequestAbortedError,
- InvalidArgumentError,
- HeadersTimeoutError,
- HeadersOverflowError,
- SocketError,
- InformationalError,
- BodyTimeoutError,
- HTTPParserError,
- ResponseExceededMaxSizeError
- } = require('../core/errors.js')
- const {
- kUrl,
- kReset,
- kClient,
- kParser,
- kBlocking,
- kRunning,
- kPending,
- kSize,
- kWriting,
- kQueue,
- kNoRef,
- kKeepAliveDefaultTimeout,
- kHostHeader,
- kPendingIdx,
- kRunningIdx,
- kError,
- kPipelining,
- kSocket,
- kKeepAliveTimeoutValue,
- kMaxHeadersSize,
- kKeepAliveMaxTimeout,
- kKeepAliveTimeoutThreshold,
- kHeadersTimeout,
- kBodyTimeout,
- kStrictContentLength,
- kMaxRequests,
- kCounter,
- kMaxResponseSize,
- kOnError,
- kResume,
- kHTTPContext,
- kClosed
- } = require('../core/symbols.js')
- const constants = require('../llhttp/constants.js')
- const EMPTY_BUF = Buffer.alloc(0)
- const FastBuffer = Buffer[Symbol.species]
- const removeAllListeners = util.removeAllListeners
- const kIdleSocketValidation = Symbol('kIdleSocketValidation')
- const kIdleSocketValidationTimeout = Symbol('kIdleSocketValidationTimeout')
- const kSocketUsed = Symbol('kSocketUsed')
- let extractBody
- function lazyllhttp () {
- const llhttpWasmData = process.env.JEST_WORKER_ID ? require('../llhttp/llhttp-wasm.js') : undefined
- let mod
- // We disable wasm SIMD on older versions of Node.js on ppc64 that are broken on Power >=9 architectures.
- let useWasmSIMD = true
- if (process.arch === 'ppc64') {
- const [major, minor] = process.versions.node.split('.').map(n => parseInt(n, 10))
- if (major < 24 || (major === 24 && minor < 12)) {
- useWasmSIMD = false
- }
- }
- // The Env Variable UNDICI_NO_WASM_SIMD allows explicitly overriding the default behavior
- if (process.env.UNDICI_NO_WASM_SIMD === '1') {
- useWasmSIMD = false
- } else if (process.env.UNDICI_NO_WASM_SIMD === '0') {
- useWasmSIMD = true
- }
- if (useWasmSIMD) {
- try {
- mod = new WebAssembly.Module(require('../llhttp/llhttp_simd-wasm.js'))
- } catch {
- }
- }
- if (!mod) {
- // We could check if the error was caused by the simd option not
- // being enabled, but the occurring of this other error
- // * https://github.com/emscripten-core/emscripten/issues/11495
- // got me to remove that check to avoid breaking Node 12.
- mod = new WebAssembly.Module(llhttpWasmData || require('../llhttp/llhttp-wasm.js'))
- }
- return new WebAssembly.Instance(mod, {
- env: {
- /**
- * @param {number} p
- * @param {number} at
- * @param {number} len
- * @returns {number}
- */
- wasm_on_url: (p, at, len) => {
- return 0
- },
- /**
- * @param {number} p
- * @param {number} at
- * @param {number} len
- * @returns {number}
- */
- wasm_on_status: (p, at, len) => {
- assert(currentParser.ptr === p)
- const start = at - currentBufferPtr + currentBufferRef.byteOffset
- return currentParser.onStatus(new FastBuffer(currentBufferRef.buffer, start, len))
- },
- /**
- * @param {number} p
- * @returns {number}
- */
- wasm_on_message_begin: (p) => {
- assert(currentParser.ptr === p)
- return currentParser.onMessageBegin()
- },
- /**
- * @param {number} p
- * @param {number} at
- * @param {number} len
- * @returns {number}
- */
- wasm_on_header_field: (p, at, len) => {
- assert(currentParser.ptr === p)
- const start = at - currentBufferPtr + currentBufferRef.byteOffset
- return currentParser.onHeaderField(new FastBuffer(currentBufferRef.buffer, start, len))
- },
- /**
- * @param {number} p
- * @param {number} at
- * @param {number} len
- * @returns {number}
- */
- wasm_on_header_value: (p, at, len) => {
- assert(currentParser.ptr === p)
- const start = at - currentBufferPtr + currentBufferRef.byteOffset
- return currentParser.onHeaderValue(new FastBuffer(currentBufferRef.buffer, start, len))
- },
- /**
- * @param {number} p
- * @param {number} statusCode
- * @param {0|1} upgrade
- * @param {0|1} shouldKeepAlive
- * @returns {number}
- */
- wasm_on_headers_complete: (p, statusCode, upgrade, shouldKeepAlive) => {
- assert(currentParser.ptr === p)
- return currentParser.onHeadersComplete(statusCode, upgrade === 1, shouldKeepAlive === 1)
- },
- /**
- * @param {number} p
- * @param {number} at
- * @param {number} len
- * @returns {number}
- */
- wasm_on_body: (p, at, len) => {
- assert(currentParser.ptr === p)
- const start = at - currentBufferPtr + currentBufferRef.byteOffset
- return currentParser.onBody(new FastBuffer(currentBufferRef.buffer, start, len))
- },
- /**
- * @param {number} p
- * @returns {number}
- */
- wasm_on_message_complete: (p) => {
- assert(currentParser.ptr === p)
- return currentParser.onMessageComplete()
- }
- }
- })
- }
- let llhttpInstance = null
- /**
- * @type {Parser|null}
- */
- let currentParser = null
- let currentBufferRef = null
- /**
- * @type {number}
- */
- let currentBufferSize = 0
- let currentBufferPtr = null
- const USE_NATIVE_TIMER = 0
- const USE_FAST_TIMER = 1
- // Use fast timers for headers and body to take eventual event loop
- // latency into account.
- const TIMEOUT_HEADERS = 2 | USE_FAST_TIMER
- const TIMEOUT_BODY = 4 | USE_FAST_TIMER
- // Use native timers to ignore event loop latency for keep-alive
- // handling.
- const TIMEOUT_KEEP_ALIVE = 8 | USE_NATIVE_TIMER
- class Parser {
- /**
- * @param {import('./client.js')} client
- * @param {import('net').Socket} socket
- * @param {*} llhttp
- */
- constructor (client, socket, { exports }) {
- this.llhttp = exports
- this.ptr = this.llhttp.llhttp_alloc(constants.TYPE.RESPONSE)
- this.client = client
- /**
- * @type {import('net').Socket}
- */
- this.socket = socket
- this.timeout = null
- this.timeoutWeakRef = new WeakRef(this)
- this.timeoutValue = null
- this.timeoutType = null
- this.statusCode = 0
- this.statusText = ''
- this.upgrade = false
- this.headers = []
- this.headersSize = 0
- this.headersMaxSize = client[kMaxHeadersSize]
- this.shouldKeepAlive = false
- this.paused = false
- this.resume = this.resume.bind(this)
- this.bytesRead = 0
- this.keepAlive = ''
- this.contentLength = ''
- this.connection = ''
- this.maxResponseSize = client[kMaxResponseSize]
- }
- setTimeout (delay, type) {
- // If the existing timer and the new timer are of different timer type
- // (fast or native) or have different delay, we need to clear the existing
- // timer and set a new one.
- if (
- delay !== this.timeoutValue ||
- (type & USE_FAST_TIMER) ^ (this.timeoutType & USE_FAST_TIMER)
- ) {
- // If a timeout is already set, clear it with clearTimeout of the fast
- // timer implementation, as it can clear fast and native timers.
- if (this.timeout) {
- timers.clearTimeout(this.timeout)
- this.timeout = null
- }
- if (delay) {
- if (type & USE_FAST_TIMER) {
- this.timeout = timers.setFastTimeout(onParserTimeout, delay, this.timeoutWeakRef)
- } else {
- this.timeout = setTimeout(onParserTimeout, delay, this.timeoutWeakRef)
- this.timeout?.unref()
- }
- }
- this.timeoutValue = delay
- } else if (this.timeout) {
- if (this.timeout.refresh) {
- this.timeout.refresh()
- }
- }
- this.timeoutType = type
- }
- resume () {
- if (this.socket.destroyed || !this.paused) {
- return
- }
- assert(this.ptr != null)
- assert(currentParser === null)
- this.llhttp.llhttp_resume(this.ptr)
- assert(this.timeoutType === TIMEOUT_BODY)
- if (this.timeout) {
- if (this.timeout.refresh) {
- this.timeout.refresh()
- }
- }
- this.paused = false
- this.execute(this.socket.read() || EMPTY_BUF) // Flush parser.
- this.readMore()
- }
- readMore () {
- while (!this.paused && this.ptr) {
- const chunk = this.socket.read()
- if (chunk === null) {
- break
- }
- this.execute(chunk)
- }
- }
- /**
- * @param {Buffer} chunk
- */
- execute (chunk) {
- assert(currentParser === null)
- assert(this.ptr != null)
- assert(!this.paused)
- const { socket, llhttp } = this
- // Allocate a new buffer if the current buffer is too small.
- if (chunk.length > currentBufferSize) {
- if (currentBufferPtr) {
- llhttp.free(currentBufferPtr)
- }
- // Allocate a buffer that is a multiple of 4096 bytes.
- currentBufferSize = Math.ceil(chunk.length / 4096) * 4096
- currentBufferPtr = llhttp.malloc(currentBufferSize)
- }
- new Uint8Array(llhttp.memory.buffer, currentBufferPtr, currentBufferSize).set(chunk)
- // Call `execute` on the wasm parser.
- // We pass the `llhttp_parser` pointer address, the pointer address of buffer view data,
- // and finally the length of bytes to parse.
- // The return value is an error code or `constants.ERROR.OK`.
- try {
- let ret
- try {
- currentBufferRef = chunk
- currentParser = this
- ret = llhttp.llhttp_execute(this.ptr, currentBufferPtr, chunk.length)
- } finally {
- currentParser = null
- currentBufferRef = null
- }
- if (ret !== constants.ERROR.OK) {
- const data = chunk.subarray(llhttp.llhttp_get_error_pos(this.ptr) - currentBufferPtr)
- if (ret === constants.ERROR.PAUSED_UPGRADE) {
- this.onUpgrade(data)
- } else if (ret === constants.ERROR.PAUSED) {
- this.paused = true
- socket.unshift(data)
- } else {
- throw this.createError(ret, data)
- }
- }
- } catch (err) {
- util.destroy(socket, err)
- }
- }
- finish () {
- assert(currentParser === null)
- assert(this.ptr != null)
- assert(!this.paused)
- const { llhttp } = this
- let ret
- try {
- currentParser = this
- ret = llhttp.llhttp_finish(this.ptr)
- } finally {
- currentParser = null
- }
- if (ret === constants.ERROR.OK) {
- return null
- }
- if (ret === constants.ERROR.PAUSED || ret === constants.ERROR.PAUSED_UPGRADE) {
- this.paused = true
- return null
- }
- return this.createError(ret, EMPTY_BUF)
- }
- createError (ret, data) {
- const { llhttp, contentLength, bytesRead } = this
- if (contentLength && bytesRead !== parseInt(contentLength, 10)) {
- return new ResponseContentLengthMismatchError()
- }
- const ptr = llhttp.llhttp_get_error_reason(this.ptr)
- let message = ''
- if (ptr) {
- const len = new Uint8Array(llhttp.memory.buffer, ptr).indexOf(0)
- message =
- 'Response does not match the HTTP/1.1 protocol (' +
- Buffer.from(llhttp.memory.buffer, ptr, len).toString() +
- ')'
- }
- return new HTTPParserError(message, constants.ERROR[ret], data)
- }
- destroy () {
- assert(currentParser === null)
- assert(this.ptr != null)
- this.llhttp.llhttp_free(this.ptr)
- this.ptr = null
- this.timeout && timers.clearTimeout(this.timeout)
- this.timeout = null
- this.timeoutValue = null
- this.timeoutType = null
- this.paused = false
- }
- /**
- * @param {Buffer} buf
- * @returns {0}
- */
- onStatus (buf) {
- this.statusText = buf.toString()
- return 0
- }
- /**
- * @returns {0|-1}
- */
- onMessageBegin () {
- const { socket, client } = this
- if (socket.destroyed) {
- return -1
- }
- if (client[kRunning] === 0) {
- util.destroy(socket, new SocketError('bad response', util.getSocketInfo(socket)))
- return -1
- }
- const request = client[kQueue][client[kRunningIdx]]
- if (!request) {
- return -1
- }
- request.onResponseStarted()
- return 0
- }
- /**
- * @param {Buffer} buf
- * @returns {number}
- */
- onHeaderField (buf) {
- const len = this.headers.length
- if ((len & 1) === 0) {
- this.headers.push(buf)
- } else {
- this.headers[len - 1] = Buffer.concat([this.headers[len - 1], buf])
- }
- this.trackHeader(buf.length)
- return 0
- }
- /**
- * @param {Buffer} buf
- * @returns {number}
- */
- onHeaderValue (buf) {
- let len = this.headers.length
- if ((len & 1) === 1) {
- this.headers.push(buf)
- len += 1
- } else {
- this.headers[len - 1] = Buffer.concat([this.headers[len - 1], buf])
- }
- const key = this.headers[len - 2]
- if (key.length === 10) {
- const headerName = util.bufferToLowerCasedHeaderName(key)
- if (headerName === 'keep-alive') {
- this.keepAlive += buf.toString()
- } else if (headerName === 'connection') {
- this.connection += buf.toString()
- }
- } else if (key.length === 14 && util.bufferToLowerCasedHeaderName(key) === 'content-length') {
- this.contentLength += buf.toString()
- }
- this.trackHeader(buf.length)
- return 0
- }
- /**
- * @param {number} len
- */
- trackHeader (len) {
- this.headersSize += len
- if (this.headersSize >= this.headersMaxSize) {
- util.destroy(this.socket, new HeadersOverflowError())
- }
- }
- /**
- * @param {Buffer} head
- */
- onUpgrade (head) {
- const { upgrade, client, socket, headers, statusCode, statusText } = this
- assert(upgrade)
- assert(client[kSocket] === socket)
- assert(!socket.destroyed)
- assert(!this.paused)
- assert((headers.length & 1) === 0)
- const request = client[kQueue][client[kRunningIdx]]
- assert(request)
- assert(request.upgrade || request.method === 'CONNECT')
- this.statusCode = 0
- this.statusText = ''
- this.shouldKeepAlive = false
- this.headers = []
- this.headersSize = 0
- socket.unshift(head)
- socket[kParser].destroy()
- socket[kParser] = null
- socket[kClient] = null
- socket[kError] = null
- removeAllListeners(socket)
- client[kSocket] = null
- client[kHTTPContext] = null // TODO (fix): This is hacky...
- client[kQueue][client[kRunningIdx]++] = null
- client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade'))
- try {
- request.onUpgrade(statusCode, headers, socket, statusText)
- } catch (err) {
- util.errorRequest(client, request, err)
- util.destroy(socket, err)
- }
- client[kResume]()
- }
- /**
- * @param {number} statusCode
- * @param {boolean} upgrade
- * @param {boolean} shouldKeepAlive
- * @returns {number}
- */
- onHeadersComplete (statusCode, upgrade, shouldKeepAlive) {
- const { client, socket, headers, statusText } = this
- if (socket.destroyed) {
- return -1
- }
- if (client[kRunning] === 0) {
- util.destroy(socket, new SocketError('bad response', util.getSocketInfo(socket)))
- return -1
- }
- const request = client[kQueue][client[kRunningIdx]]
- if (!request) {
- return -1
- }
- assert(!this.upgrade)
- assert(this.statusCode < 200)
- if (statusCode === 100) {
- util.destroy(socket, new SocketError('bad response', util.getSocketInfo(socket)))
- return -1
- }
- /* this can only happen if server is misbehaving */
- if (upgrade && !request.upgrade) {
- util.destroy(socket, new SocketError('bad upgrade', util.getSocketInfo(socket)))
- return -1
- }
- assert(this.timeoutType === TIMEOUT_HEADERS)
- this.statusCode = statusCode
- this.shouldKeepAlive = (
- shouldKeepAlive ||
- // Override llhttp value which does not allow keepAlive for HEAD.
- (request.method === 'HEAD' && !socket[kReset] && this.connection.toLowerCase() === 'keep-alive')
- )
- if (this.statusCode >= 200) {
- const bodyTimeout = request.bodyTimeout != null
- ? request.bodyTimeout
- : client[kBodyTimeout]
- this.setTimeout(bodyTimeout, TIMEOUT_BODY)
- } else if (this.timeout) {
- if (this.timeout.refresh) {
- this.timeout.refresh()
- }
- }
- if (request.method === 'CONNECT') {
- assert(client[kRunning] === 1)
- this.upgrade = true
- return 2
- }
- if (upgrade) {
- assert(client[kRunning] === 1)
- this.upgrade = true
- return 2
- }
- assert((this.headers.length & 1) === 0)
- this.headers = []
- this.headersSize = 0
- if (this.shouldKeepAlive && client[kPipelining]) {
- const keepAliveTimeout = this.keepAlive ? util.parseKeepAliveTimeout(this.keepAlive) : null
- if (keepAliveTimeout != null) {
- const timeout = Math.min(
- keepAliveTimeout - client[kKeepAliveTimeoutThreshold],
- client[kKeepAliveMaxTimeout]
- )
- if (timeout <= 0) {
- socket[kReset] = true
- } else {
- client[kKeepAliveTimeoutValue] = timeout
- }
- } else {
- client[kKeepAliveTimeoutValue] = client[kKeepAliveDefaultTimeout]
- }
- } else {
- // Stop more requests from being dispatched.
- socket[kReset] = true
- }
- const pause = request.onHeaders(statusCode, headers, this.resume, statusText) === false
- if (request.aborted) {
- return -1
- }
- if (request.method === 'HEAD') {
- return 1
- }
- if (statusCode < 200) {
- return 1
- }
- if (socket[kBlocking]) {
- socket[kBlocking] = false
- client[kResume]()
- }
- return pause ? constants.ERROR.PAUSED : 0
- }
- /**
- * @param {Buffer} buf
- * @returns {number}
- */
- onBody (buf) {
- const { client, socket, statusCode, maxResponseSize } = this
- if (socket.destroyed) {
- return -1
- }
- const request = client[kQueue][client[kRunningIdx]]
- assert(request)
- assert(this.timeoutType === TIMEOUT_BODY)
- if (this.timeout) {
- if (this.timeout.refresh) {
- this.timeout.refresh()
- }
- }
- assert(statusCode >= 200)
- if (maxResponseSize > -1 && this.bytesRead + buf.length > maxResponseSize) {
- util.destroy(socket, new ResponseExceededMaxSizeError())
- return -1
- }
- this.bytesRead += buf.length
- if (request.onData(buf) === false) {
- return constants.ERROR.PAUSED
- }
- return 0
- }
- /**
- * @returns {number}
- */
- onMessageComplete () {
- const { client, socket, statusCode, upgrade, headers, contentLength, bytesRead, shouldKeepAlive } = this
- if (socket.destroyed && (!statusCode || shouldKeepAlive)) {
- return -1
- }
- if (upgrade) {
- return 0
- }
- assert(statusCode >= 100)
- assert((this.headers.length & 1) === 0)
- const request = client[kQueue][client[kRunningIdx]]
- assert(request)
- this.statusCode = 0
- this.statusText = ''
- this.bytesRead = 0
- this.contentLength = ''
- this.keepAlive = ''
- this.connection = ''
- this.headers = []
- this.headersSize = 0
- if (statusCode < 200) {
- return 0
- }
- if (request.method !== 'HEAD' && contentLength && bytesRead !== parseInt(contentLength, 10)) {
- util.destroy(socket, new ResponseContentLengthMismatchError())
- return -1
- }
- request.onComplete(headers)
- client[kQueue][client[kRunningIdx]++] = null
- socket[kSocketUsed] = client[kPending] === 0
- if (socket[kWriting]) {
- assert(client[kRunning] === 0)
- // Response completed before request.
- util.destroy(socket, new InformationalError('reset'))
- return constants.ERROR.PAUSED
- } else if (!shouldKeepAlive) {
- util.destroy(socket, new InformationalError('reset'))
- return constants.ERROR.PAUSED
- } else if (socket[kReset] && client[kRunning] === 0) {
- // Destroy socket once all requests have completed.
- // The request at the tail of the pipeline is the one
- // that requested reset and no further requests should
- // have been queued since then.
- util.destroy(socket, new InformationalError('reset'))
- return constants.ERROR.PAUSED
- } else if (client[kPipelining] == null || client[kPipelining] === 1) {
- // We must wait a full event loop cycle to reuse this socket to make sure
- // that non-spec compliant servers are not closing the connection even if they
- // said they won't.
- setImmediate(client[kResume])
- } else {
- client[kResume]()
- }
- return 0
- }
- }
- function onParserTimeout (parserWeakRef) {
- const parser = parserWeakRef.deref()
- if (!parser) {
- return
- }
- const { socket, timeoutType, client, paused } = parser
- if (timeoutType === TIMEOUT_HEADERS) {
- if (!socket[kWriting] || socket.writableNeedDrain || client[kRunning] > 1) {
- assert(!paused, 'cannot be paused while waiting for headers')
- util.destroy(socket, new HeadersTimeoutError())
- }
- } else if (timeoutType === TIMEOUT_BODY) {
- if (!paused) {
- util.destroy(socket, new BodyTimeoutError())
- }
- } else if (timeoutType === TIMEOUT_KEEP_ALIVE) {
- assert(client[kRunning] === 0 && client[kKeepAliveTimeoutValue])
- util.destroy(socket, new InformationalError('socket idle timeout'))
- }
- }
- /**
- * @param {import ('./client.js')} client
- * @param {import('net').Socket} socket
- * @returns
- */
- function connectH1 (client, socket) {
- client[kSocket] = socket
- if (!llhttpInstance) {
- llhttpInstance = lazyllhttp()
- }
- if (socket.errored) {
- throw socket.errored
- }
- if (socket.destroyed) {
- throw new SocketError('destroyed')
- }
- socket[kNoRef] = false
- socket[kWriting] = false
- socket[kReset] = false
- socket[kBlocking] = false
- socket[kIdleSocketValidation] = 0
- socket[kIdleSocketValidationTimeout] = null
- socket[kSocketUsed] = false
- socket[kParser] = new Parser(client, socket, llhttpInstance)
- util.addListener(socket, 'error', onHttpSocketError)
- util.addListener(socket, 'readable', onHttpSocketReadable)
- util.addListener(socket, 'end', onHttpSocketEnd)
- util.addListener(socket, 'close', onHttpSocketClose)
- socket[kClosed] = false
- socket.on('close', onSocketClose)
- return {
- version: 'h1',
- defaultPipelining: 1,
- write (request) {
- return writeH1(client, request)
- },
- resume () {
- resumeH1(client)
- },
- /**
- * @param {Error|undefined} err
- * @param {() => void} callback
- */
- destroy (err, callback) {
- if (socket[kClosed]) {
- queueMicrotask(callback)
- } else {
- socket.on('close', callback)
- socket.destroy(err)
- }
- },
- /**
- * @returns {boolean}
- */
- get destroyed () {
- return socket.destroyed
- },
- /**
- * @param {import('../core/request.js')} request
- * @returns {boolean}
- */
- busy (request) {
- if (socket[kWriting] || socket[kReset] || socket[kBlocking] || socket[kIdleSocketValidation] === 1) {
- return true
- }
- if (request) {
- if (client[kRunning] > 0 && !request.idempotent) {
- // Non-idempotent request cannot be retried.
- // Ensure that no other requests are inflight and
- // could cause failure.
- return true
- }
- if (client[kRunning] > 0 && (request.upgrade || request.method === 'CONNECT')) {
- // Don't dispatch an upgrade until all preceding requests have completed.
- // A misbehaving server might upgrade the connection before all pipelined
- // request has completed.
- return true
- }
- if (client[kRunning] > 0 && util.bodyLength(request.body) !== 0 &&
- (util.isStream(request.body) || util.isAsyncIterable(request.body) || util.isFormDataLike(request.body))) {
- // Request with stream or iterator body can error while other requests
- // are inflight and indirectly error those as well.
- // Ensure this doesn't happen by waiting for inflight
- // to complete before dispatching.
- // Request with stream or iterator body cannot be retried.
- // Ensure that no other requests are inflight and
- // could cause failure.
- return true
- }
- }
- return false
- }
- }
- }
- function onHttpSocketError (err) {
- assert(err.code !== 'ERR_TLS_CERT_ALTNAME_INVALID')
- const parser = this[kParser]
- // On Mac OS, we get an ECONNRESET even if there is a full body to be forwarded
- // to the user.
- if (err.code === 'ECONNRESET' && parser.statusCode && !parser.shouldKeepAlive) {
- const parserErr = parser.finish()
- if (parserErr) {
- this[kError] = parserErr
- this[kClient][kOnError](parserErr)
- }
- return
- }
- this[kError] = err
- this[kClient][kOnError](err)
- }
- function onHttpSocketReadable () {
- this[kParser]?.readMore()
- }
- function onHttpSocketEnd () {
- const parser = this[kParser]
- if (parser.statusCode && !parser.shouldKeepAlive) {
- const parserErr = parser.finish()
- if (parserErr) {
- util.destroy(this, parserErr)
- }
- return
- }
- util.destroy(this, new SocketError('other side closed', util.getSocketInfo(this)))
- }
- function onHttpSocketClose () {
- const parser = this[kParser]
- clearIdleSocketValidation(this)
- if (parser) {
- if (!this[kError] && parser.statusCode && !parser.shouldKeepAlive) {
- this[kError] = parser.finish() || this[kError]
- }
- this[kParser].destroy()
- this[kParser] = null
- }
- const err = this[kError] || new SocketError('closed', util.getSocketInfo(this))
- const client = this[kClient]
- client[kSocket] = null
- client[kHTTPContext] = null // TODO (fix): This is hacky...
- if (client.destroyed) {
- assert(client[kPending] === 0)
- // Fail entire queue.
- const requests = client[kQueue].splice(client[kRunningIdx])
- for (let i = 0; i < requests.length; i++) {
- const request = requests[i]
- util.errorRequest(client, request, err)
- }
- } else if (client[kRunning] > 0 && err.code !== 'UND_ERR_INFO') {
- // Fail head of pipeline.
- const request = client[kQueue][client[kRunningIdx]]
- client[kQueue][client[kRunningIdx]++] = null
- util.errorRequest(client, request, err)
- }
- client[kPendingIdx] = client[kRunningIdx]
- assert(client[kRunning] === 0)
- client.emit('disconnect', client[kUrl], [client], err)
- client[kResume]()
- }
- function onSocketClose () {
- this[kClosed] = true
- }
- function clearIdleSocketValidation (socket) {
- if (socket[kIdleSocketValidationTimeout]) {
- clearImmediate(socket[kIdleSocketValidationTimeout])
- socket[kIdleSocketValidationTimeout] = null
- }
- socket[kIdleSocketValidation] = 0
- }
- function scheduleIdleSocketValidation (client, socket) {
- socket[kIdleSocketValidation] = 1
- // Yield to the check phase (after poll) so unsolicited bytes / FIN / RST
- // already pending on this idle keep-alive socket are processed before the
- // next request is written (GHSA-35p6-xmwp-9g52).
- //
- // setTimeout(0) pays Node's ~1ms timer floor on every sequential reuse
- // (#5493). setImmediate avoids that, but an *unref'd* Immediate lets poll
- // block for ~500ms when the event loop is otherwise idle (#5600 / #5606).
- // A ref'd Immediate both keeps the pending request alive and makes poll
- // return immediately — the hybrid those issues asked for.
- socket[kIdleSocketValidationTimeout] = setImmediate(() => {
- socket[kIdleSocketValidationTimeout] = null
- socket[kIdleSocketValidation] = 2
- if (client[kSocket] === socket && !socket.destroyed) {
- client[kResume]()
- }
- })
- }
- /**
- * @param {import('./client.js')} client
- */
- function resumeH1 (client) {
- const socket = client[kSocket]
- if (socket && !socket.destroyed) {
- if (client[kSize] === 0) {
- if (!socket[kNoRef] && socket.unref) {
- socket.unref()
- socket[kNoRef] = true
- }
- } else if (socket[kNoRef] && socket.ref) {
- socket.ref()
- socket[kNoRef] = false
- }
- if (client[kRunning] === 0 && client[kPending] > 0 && socket[kSocketUsed]) {
- if (socket[kIdleSocketValidation] === 0) {
- scheduleIdleSocketValidation(client, socket)
- socket[kParser].readMore()
- if (socket.destroyed) {
- return
- }
- return
- }
- if (socket[kIdleSocketValidation] === 1) {
- socket[kParser].readMore()
- if (socket.destroyed) {
- return
- }
- return
- }
- }
- if (client[kRunning] === 0) {
- socket[kParser].readMore()
- if (socket.destroyed) {
- return
- }
- }
- if (client[kSize] === 0) {
- if (socket[kParser].timeoutType !== TIMEOUT_KEEP_ALIVE) {
- socket[kParser].setTimeout(client[kKeepAliveTimeoutValue], TIMEOUT_KEEP_ALIVE)
- }
- } else if (client[kRunning] > 0 && socket[kParser].statusCode < 200) {
- if (socket[kParser].timeoutType !== TIMEOUT_HEADERS) {
- const request = client[kQueue][client[kRunningIdx]]
- const headersTimeout = request.headersTimeout != null
- ? request.headersTimeout
- : client[kHeadersTimeout]
- socket[kParser].setTimeout(headersTimeout, TIMEOUT_HEADERS)
- }
- }
- }
- }
- // https://www.rfc-editor.org/rfc/rfc7230#section-3.3.2
- function shouldSendContentLength (method) {
- return method !== 'GET' && method !== 'HEAD' && method !== 'OPTIONS' && method !== 'TRACE' && method !== 'CONNECT'
- }
- /**
- * @param {import('./client.js')} client
- * @param {import('../core/request.js')} request
- * @returns
- */
- function writeH1 (client, request) {
- const { method, path, host, upgrade, blocking, reset } = request
- let { body, headers, contentLength } = request
- // https://tools.ietf.org/html/rfc7231#section-4.3.1
- // https://tools.ietf.org/html/rfc7231#section-4.3.2
- // https://tools.ietf.org/html/rfc7231#section-4.3.5
- // Sending a payload body on a request that does not
- // expect it can cause undefined behavior on some
- // servers and corrupt connection state. Do not
- // re-use the connection for further requests.
- const expectsPayload = (
- method === 'PUT' ||
- method === 'POST' ||
- method === 'PATCH' ||
- method === 'QUERY' ||
- method === 'PROPFIND' ||
- method === 'PROPPATCH'
- )
- if (util.isFormDataLike(body)) {
- if (!extractBody) {
- extractBody = require('../web/fetch/body.js').extractBody
- }
- const [bodyStream, contentType] = extractBody(body)
- if (request.contentType == null) {
- headers.push('content-type', contentType)
- }
- body = bodyStream.stream
- contentLength = bodyStream.length
- } else if (util.isBlobLike(body) && request.contentType == null) {
- const contentType = body.type
- if (contentType) {
- const contentTypeValue = `${contentType}`
- if (!util.isValidHeaderValue(contentTypeValue)) {
- util.errorRequest(client, request, new InvalidArgumentError('invalid content-type header'))
- return false
- }
- headers.push('content-type', contentTypeValue)
- }
- }
- if (body && typeof body.read === 'function') {
- // Try to read EOF in order to get length.
- body.read(0)
- }
- const bodyLength = util.bodyLength(body)
- contentLength = bodyLength ?? contentLength
- if (contentLength === null) {
- contentLength = request.contentLength
- }
- if (contentLength === 0 && !expectsPayload) {
- // https://tools.ietf.org/html/rfc7230#section-3.3.2
- // A user agent SHOULD NOT send a Content-Length header field when
- // the request message does not contain a payload body and the method
- // semantics do not anticipate such a body.
- contentLength = null
- }
- // https://github.com/nodejs/undici/issues/2046
- // A user agent may send a Content-Length header with 0 value, this should be allowed.
- if (shouldSendContentLength(method) && contentLength > 0 && request.contentLength !== null && request.contentLength !== contentLength) {
- if (client[kStrictContentLength]) {
- util.errorRequest(client, request, new RequestContentLengthMismatchError())
- return false
- }
- process.emitWarning(new RequestContentLengthMismatchError())
- }
- const socket = client[kSocket]
- clearIdleSocketValidation(socket)
- /**
- * @param {Error} [err]
- * @returns {void}
- */
- const abort = (err) => {
- if (request.aborted || request.completed) {
- return
- }
- util.errorRequest(client, request, err || new RequestAbortedError())
- util.destroy(body)
- util.destroy(socket, new InformationalError('aborted'))
- }
- try {
- request.onConnect(abort)
- } catch (err) {
- util.errorRequest(client, request, err)
- }
- if (request.aborted) {
- return false
- }
- if (method === 'HEAD') {
- // https://github.com/mcollina/undici/issues/258
- // Close after a HEAD request to interop with misbehaving servers
- // that may send a body in the response.
- socket[kReset] = true
- }
- if (upgrade || method === 'CONNECT') {
- // On CONNECT or upgrade, block pipeline from dispatching further
- // requests on this connection.
- socket[kReset] = true
- }
- if (reset != null) {
- socket[kReset] = reset
- }
- if (client[kMaxRequests] && socket[kCounter]++ >= client[kMaxRequests]) {
- socket[kReset] = true
- }
- if (blocking) {
- socket[kBlocking] = true
- }
- if (socket.setTypeOfService) {
- socket.setTypeOfService(request.typeOfService)
- }
- let header = `${method} ${path} HTTP/1.1\r\n`
- if (typeof host === 'string') {
- header += `host: ${host}\r\n`
- } else {
- header += client[kHostHeader]
- }
- if (upgrade) {
- header += `connection: upgrade\r\nupgrade: ${upgrade}\r\n`
- } else if (client[kPipelining] && !socket[kReset]) {
- header += 'connection: keep-alive\r\n'
- } else {
- header += 'connection: close\r\n'
- }
- if (Array.isArray(headers)) {
- for (let n = 0; n < headers.length; n += 2) {
- const key = headers[n + 0]
- const val = headers[n + 1]
- if (Array.isArray(val)) {
- for (let i = 0; i < val.length; i++) {
- header += `${key}: ${val[i]}\r\n`
- }
- } else {
- header += `${key}: ${val}\r\n`
- }
- }
- }
- if (channels.sendHeaders.hasSubscribers) {
- channels.sendHeaders.publish({ request, headers: header, socket })
- }
- if (!body || bodyLength === 0) {
- writeBuffer(abort, null, client, request, socket, contentLength, header, expectsPayload)
- } else if (util.isBuffer(body)) {
- writeBuffer(abort, body, client, request, socket, contentLength, header, expectsPayload)
- } else if (util.isBlobLike(body)) {
- if (typeof body.stream === 'function') {
- writeIterable(abort, body.stream(), client, request, socket, contentLength, header, expectsPayload)
- } else {
- writeBlob(abort, body, client, request, socket, contentLength, header, expectsPayload)
- }
- } else if (util.isStream(body)) {
- writeStream(abort, body, client, request, socket, contentLength, header, expectsPayload)
- } else if (util.isIterable(body)) {
- writeIterable(abort, body, client, request, socket, contentLength, header, expectsPayload)
- } else {
- assert(false)
- }
- return true
- }
- /**
- * @param {AbortCallback} abort
- * @param {import('stream').Stream} body
- * @param {import('./client.js')} client
- * @param {import('../core/request.js')} request
- * @param {import('net').Socket} socket
- * @param {number} contentLength
- * @param {string} header
- * @param {boolean} expectsPayload
- */
- function writeStream (abort, body, client, request, socket, contentLength, header, expectsPayload) {
- assert(contentLength !== 0 || client[kRunning] === 0, 'stream body cannot be pipelined')
- let finished = false
- const writer = new AsyncWriter({ abort, socket, request, contentLength, client, expectsPayload, header })
- /**
- * @param {Buffer} chunk
- * @returns {void}
- */
- const onData = function (chunk) {
- if (finished) {
- return
- }
- try {
- if (!writer.write(chunk) && this.pause) {
- this.pause()
- }
- } catch (err) {
- util.destroy(this, err)
- }
- }
- /**
- * @returns {void}
- */
- const onDrain = function () {
- if (finished) {
- return
- }
- if (body.resume) {
- body.resume()
- }
- }
- /**
- * @returns {void}
- */
- const onClose = function () {
- // 'close' might be emitted *before* 'error' for
- // broken streams. Wait a tick to avoid this case.
- queueMicrotask(() => {
- // It's only safe to remove 'error' listener after
- // 'close'.
- body.removeListener('error', onFinished)
- })
- if (!finished) {
- const err = new RequestAbortedError()
- queueMicrotask(() => onFinished(err))
- }
- }
- /**
- * @param {Error} [err]
- * @returns
- */
- const onFinished = function (err) {
- if (finished) {
- return
- }
- finished = true
- assert(socket.destroyed || (socket[kWriting] && client[kRunning] <= 1))
- socket
- .off('drain', onDrain)
- .off('error', onFinished)
- body
- .removeListener('data', onData)
- .removeListener('end', onFinished)
- .removeListener('close', onClose)
- if (!err) {
- try {
- writer.end()
- } catch (er) {
- err = er
- }
- }
- writer.destroy(err)
- if (err && (err.code !== 'UND_ERR_INFO' || err.message !== 'reset')) {
- util.destroy(body, err)
- } else {
- util.destroy(body)
- }
- }
- body
- .on('data', onData)
- .on('end', onFinished)
- .on('error', onFinished)
- .on('close', onClose)
- if (body.resume) {
- body.resume()
- }
- socket
- .on('drain', onDrain)
- .on('error', onFinished)
- if (body.errorEmitted ?? body.errored) {
- setImmediate(onFinished, body.errored)
- } else if (body.endEmitted ?? body.readableEnded) {
- setImmediate(onFinished, null)
- }
- if (body.closeEmitted ?? body.closed) {
- setImmediate(onClose)
- }
- }
- /**
- * @typedef AbortCallback
- * @type {Function}
- * @param {Error} [err]
- * @returns {void}
- */
- /**
- * @param {AbortCallback} abort
- * @param {Uint8Array|null} body
- * @param {import('./client.js')} client
- * @param {import('../core/request.js')} request
- * @param {import('net').Socket} socket
- * @param {number} contentLength
- * @param {string} header
- * @param {boolean} expectsPayload
- * @returns {void}
- */
- function writeBuffer (abort, body, client, request, socket, contentLength, header, expectsPayload) {
- try {
- if (!body) {
- if (contentLength === 0) {
- socket.write(`${header}content-length: 0\r\n\r\n`, 'latin1')
- } else {
- assert(contentLength === null, 'no body must not have content length')
- socket.write(`${header}\r\n`, 'latin1')
- }
- } else if (util.isBuffer(body)) {
- assert(contentLength === body.byteLength, 'buffer body must have content length')
- socket.cork()
- socket.write(`${header}content-length: ${contentLength}\r\n\r\n`, 'latin1')
- socket.write(body)
- socket.uncork()
- request.onBodySent(body)
- if (!expectsPayload && request.reset !== false) {
- socket[kReset] = true
- }
- }
- request.onRequestSent()
- client[kResume]()
- } catch (err) {
- abort(err)
- }
- }
- /**
- * @param {AbortCallback} abort
- * @param {Blob} body
- * @param {import('./client.js')} client
- * @param {import('../core/request.js')} request
- * @param {import('net').Socket} socket
- * @param {number} contentLength
- * @param {string} header
- * @param {boolean} expectsPayload
- * @returns {Promise<void>}
- */
- async function writeBlob (abort, body, client, request, socket, contentLength, header, expectsPayload) {
- assert(contentLength === body.size, 'blob body must have content length')
- try {
- if (contentLength != null && contentLength !== body.size) {
- throw new RequestContentLengthMismatchError()
- }
- const buffer = Buffer.from(await body.arrayBuffer())
- socket.cork()
- socket.write(`${header}content-length: ${contentLength}\r\n\r\n`, 'latin1')
- socket.write(buffer)
- socket.uncork()
- request.onBodySent(buffer)
- request.onRequestSent()
- if (!expectsPayload && request.reset !== false) {
- socket[kReset] = true
- }
- client[kResume]()
- } catch (err) {
- abort(err)
- }
- }
- /**
- * @param {AbortCallback} abort
- * @param {Iterable} body
- * @param {import('./client.js')} client
- * @param {import('../core/request.js')} request
- * @param {import('net').Socket} socket
- * @param {number} contentLength
- * @param {string} header
- * @param {boolean} expectsPayload
- * @returns {Promise<void>}
- */
- async function writeIterable (abort, body, client, request, socket, contentLength, header, expectsPayload) {
- assert(contentLength !== 0 || client[kRunning] === 0, 'iterator body cannot be pipelined')
- let callback = null
- function onDrain () {
- if (callback) {
- const cb = callback
- callback = null
- cb()
- }
- }
- const waitForDrain = () => new Promise((resolve, reject) => {
- assert(callback === null)
- if (socket[kError]) {
- reject(socket[kError])
- } else {
- callback = resolve
- }
- })
- socket
- .on('close', onDrain)
- .on('drain', onDrain)
- const writer = new AsyncWriter({ abort, socket, request, contentLength, client, expectsPayload, header })
- try {
- // It's up to the user to somehow abort the async iterable.
- for await (const chunk of body) {
- if (socket[kError]) {
- throw socket[kError]
- }
- if (!writer.write(chunk)) {
- await waitForDrain()
- }
- }
- writer.end()
- } catch (err) {
- writer.destroy(err)
- } finally {
- socket
- .off('close', onDrain)
- .off('drain', onDrain)
- }
- }
- class AsyncWriter {
- /**
- *
- * @param {object} arg
- * @param {AbortCallback} arg.abort
- * @param {import('net').Socket} arg.socket
- * @param {import('../core/request.js')} arg.request
- * @param {number} arg.contentLength
- * @param {import('./client.js')} arg.client
- * @param {boolean} arg.expectsPayload
- * @param {string} arg.header
- */
- constructor ({ abort, socket, request, contentLength, client, expectsPayload, header }) {
- this.socket = socket
- this.request = request
- this.contentLength = contentLength
- this.client = client
- this.bytesWritten = 0
- this.expectsPayload = expectsPayload
- this.header = header
- this.abort = abort
- socket[kWriting] = true
- }
- /**
- * @param {Buffer} chunk
- * @returns
- */
- write (chunk) {
- const { socket, request, contentLength, client, bytesWritten, expectsPayload, header } = this
- if (socket[kError]) {
- throw socket[kError]
- }
- if (socket.destroyed) {
- return false
- }
- const len = Buffer.byteLength(chunk)
- if (!len) {
- return true
- }
- // We should defer writing chunks.
- if (contentLength !== null && bytesWritten + len > contentLength) {
- if (client[kStrictContentLength]) {
- throw new RequestContentLengthMismatchError()
- }
- process.emitWarning(new RequestContentLengthMismatchError())
- }
- socket.cork()
- if (bytesWritten === 0) {
- if (!expectsPayload && request.reset !== false) {
- socket[kReset] = true
- }
- if (contentLength === null) {
- socket.write(`${header}transfer-encoding: chunked\r\n`, 'latin1')
- } else {
- socket.write(`${header}content-length: ${contentLength}\r\n\r\n`, 'latin1')
- }
- }
- if (contentLength === null) {
- socket.write(`\r\n${len.toString(16)}\r\n`, 'latin1')
- }
- this.bytesWritten += len
- const ret = socket.write(chunk)
- socket.uncork()
- request.onBodySent(chunk)
- if (!ret) {
- if (socket[kParser].timeout && socket[kParser].timeoutType === TIMEOUT_HEADERS) {
- if (socket[kParser].timeout.refresh) {
- socket[kParser].timeout.refresh()
- }
- }
- }
- return ret
- }
- /**
- * @returns {void}
- */
- end () {
- const { socket, contentLength, client, bytesWritten, expectsPayload, header, request } = this
- request.onRequestSent()
- socket[kWriting] = false
- if (socket[kError]) {
- throw socket[kError]
- }
- if (socket.destroyed) {
- return
- }
- if (bytesWritten === 0) {
- if (expectsPayload) {
- // https://tools.ietf.org/html/rfc7230#section-3.3.2
- // A user agent SHOULD send a Content-Length in a request message when
- // no Transfer-Encoding is sent and the request method defines a meaning
- // for an enclosed payload body.
- socket.write(`${header}content-length: 0\r\n\r\n`, 'latin1')
- } else {
- socket.write(`${header}\r\n`, 'latin1')
- }
- } else if (contentLength === null) {
- socket.write('\r\n0\r\n\r\n', 'latin1')
- }
- if (contentLength !== null && bytesWritten !== contentLength) {
- if (client[kStrictContentLength]) {
- throw new RequestContentLengthMismatchError()
- } else {
- process.emitWarning(new RequestContentLengthMismatchError())
- }
- }
- if (socket[kParser].timeout && socket[kParser].timeoutType === TIMEOUT_HEADERS) {
- if (socket[kParser].timeout.refresh) {
- socket[kParser].timeout.refresh()
- }
- }
- client[kResume]()
- }
- /**
- * @param {Error} [err]
- * @returns {void}
- */
- destroy (err) {
- const { socket, client, abort } = this
- socket[kWriting] = false
- if (err) {
- assert(client[kRunning] <= 1, 'pipeline should only contain this request')
- abort(err)
- }
- }
- }
- module.exports = connectH1
|