| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504 |
- 'use strict'
- const assert = require('node:assert')
- const { kRetryHandlerDefaultRetry } = require('../core/symbols')
- const { RequestRetryError } = require('../core/errors')
- const WrapHandler = require('./wrap-handler')
- const {
- isDisturbed,
- parseRangeHeader,
- wrapRequestBody
- } = require('../core/util')
- function calculateRetryAfterHeader (retryAfter) {
- const retryTime = new Date(retryAfter).getTime()
- return isNaN(retryTime) ? 0 : retryTime - Date.now()
- }
- function validatePartialResponseContentLength (headers, range, statusCode, retryCount) {
- const contentLength = headers['content-length']
- if (contentLength == null) {
- return
- }
- if (!Number.isFinite(range.start) || !Number.isFinite(range.end)) {
- return
- }
- const length = Number(contentLength)
- const expectedLength = range.end - range.start + 1
- if (!Number.isFinite(length) || length !== expectedLength) {
- throw new RequestRetryError('Content-Length mismatch', statusCode, {
- headers,
- data: { count: retryCount }
- })
- }
- }
- // A stable controller handed to the downstream handler for the lifetime of the
- // request. Each transparent retry/resume is a separate dispatch with its own
- // connection controller. The proxy always forwards to the active connection
- // while preserving a downstream pause across controller replacement.
- class RetryController {
- #paused = false
- #target = null
- set target (target) {
- this.#target = target
- if (this.#paused) {
- target?.pause()
- }
- }
- get target () { return this.#target }
- pause () {
- this.#paused = true
- this.#target?.pause()
- }
- resume () {
- this.#paused = false
- this.#target?.resume()
- }
- abort (reason) {
- this.#target?.abort(reason)
- }
- get paused () { return this.#paused || (this.#target?.paused ?? false) }
- get aborted () { return this.#target?.aborted ?? false }
- get reason () { return this.#target?.reason ?? null }
- get rawHeaders () { return this.#target?.rawHeaders ?? null }
- set rawHeaders (value) {
- if (this.#target) {
- this.#target.rawHeaders = value
- }
- }
- get rawTrailers () { return this.#target?.rawTrailers ?? null }
- set rawTrailers (value) {
- if (this.#target) {
- this.#target.rawTrailers = value
- }
- }
- }
- class RetryHandler {
- constructor (opts, { dispatch, handler }) {
- const { retryOptions, ...dispatchOpts } = opts
- const {
- // Retry scoped
- retry: retryFn,
- maxRetries,
- maxTimeout,
- minTimeout,
- timeoutFactor,
- // Response scoped
- methods,
- errorCodes,
- retryAfter,
- statusCodes,
- throwOnError
- } = retryOptions ?? {}
- this.error = null
- this.dispatch = dispatch
- this.handler = WrapHandler.wrap(handler)
- this.opts = { ...dispatchOpts, body: wrapRequestBody(opts.body) }
- this.retryOpts = {
- throwOnError: throwOnError ?? true,
- retry: retryFn ?? RetryHandler[kRetryHandlerDefaultRetry],
- retryAfter: retryAfter ?? true,
- maxTimeout: maxTimeout ?? 30 * 1000, // 30s,
- minTimeout: minTimeout ?? 500, // .5s
- timeoutFactor: timeoutFactor ?? 2,
- maxRetries: maxRetries ?? 5,
- // What errors we should retry
- methods: methods ?? ['GET', 'HEAD', 'OPTIONS', 'PUT', 'DELETE', 'TRACE'],
- // Indicates which errors to retry
- statusCodes: statusCodes ?? [500, 502, 503, 504, 429],
- // List of errors to retry
- errorCodes: errorCodes ?? [
- 'ECONNRESET',
- 'ECONNREFUSED',
- 'ENOTFOUND',
- 'ENETDOWN',
- 'ENETUNREACH',
- 'EHOSTDOWN',
- 'EHOSTUNREACH',
- 'EPIPE',
- 'UND_ERR_SOCKET'
- ]
- }
- this.retryCount = 0
- this.retryCountCheckpoint = 0
- this.headersSent = false
- this.start = 0
- this.end = null
- this.etag = null
- this.controllerProxy = new RetryController()
- }
- onResponseStartWithRetry (controller, statusCode, headers, statusMessage, err) {
- if (this.retryOpts.throwOnError) {
- // Preserve old behavior for status codes that are not eligible for retry
- if (this.retryOpts.statusCodes.includes(statusCode) === false) {
- if (this.headersSent) {
- // The downstream handler already received the response from an
- // earlier attempt. Forwarding this response would replace the
- // downstream body and leave the original body pending forever.
- this.handler.onResponseError?.(this.controllerProxy, err)
- } else {
- this.headersSent = true
- this.checkpointResponseEnd(headers)
- this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage)
- }
- } else {
- this.error = err
- }
- return
- }
- if (isDisturbed(this.opts.body)) {
- this.headersSent = true
- this.checkpointResponseEnd(headers)
- this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage)
- return
- }
- function shouldRetry (passedErr) {
- if (passedErr) {
- if (this.headersSent) {
- // The downstream handler already received the response from an
- // earlier attempt. Forwarding this response would replace the
- // downstream body and leave the original body pending forever.
- this.handler.onResponseError?.(this.controllerProxy, passedErr)
- } else {
- this.headersSent = true
- this.checkpointResponseEnd(headers)
- this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage)
- }
- controller.resume()
- return
- }
- this.error = err
- controller.resume()
- }
- controller.pause()
- this.retryOpts.retry(
- err,
- {
- state: { counter: this.retryCount },
- opts: { retryOptions: this.retryOpts, ...this.opts }
- },
- shouldRetry.bind(this)
- )
- }
- checkpointResponseEnd (headers) {
- if (this.end == null && this.opts.method !== 'HEAD') {
- const contentLength = headers['content-length']
- this.end = contentLength != null ? Number(contentLength) - 1 : null
- assert(
- this.end == null || Number.isFinite(this.end),
- 'invalid content-length'
- )
- this.resume = this.end != null
- }
- }
- onRequestStart (controller, context) {
- this.controllerProxy.target = controller
- if (!this.headersSent) {
- this.handler.onRequestStart?.(this.controllerProxy, context)
- }
- }
- onRequestUpgrade (_controller, statusCode, headers, socket) {
- this.handler.onRequestUpgrade?.(this.controllerProxy, statusCode, headers, socket)
- }
- static [kRetryHandlerDefaultRetry] (err, { state, opts }, cb) {
- const { statusCode, code, headers } = err
- const { method, retryOptions } = opts
- const {
- maxRetries,
- minTimeout,
- maxTimeout,
- timeoutFactor,
- statusCodes,
- errorCodes,
- methods
- } = retryOptions
- const { counter } = state
- // Any code that is not a Undici's originated and allowed to retry
- if (code && code !== 'UND_ERR_REQ_RETRY' && !errorCodes.includes(code)) {
- cb(err)
- return
- }
- // If a set of method are provided and the current method is not in the list
- if (Array.isArray(methods) && !methods.includes(method)) {
- cb(err)
- return
- }
- // If a set of status code are provided and the current status code is not in the list
- if (
- statusCode != null &&
- Array.isArray(statusCodes) &&
- !statusCodes.includes(statusCode)
- ) {
- cb(err)
- return
- }
- // If we reached the max number of retries
- if (counter > maxRetries) {
- cb(err)
- return
- }
- let retryAfterHeader = headers?.['retry-after']
- if (retryAfterHeader) {
- retryAfterHeader = Number(retryAfterHeader)
- retryAfterHeader = Number.isNaN(retryAfterHeader)
- ? calculateRetryAfterHeader(headers['retry-after'])
- : retryAfterHeader * 1e3 // Retry-After is in seconds
- }
- const retryTimeout =
- retryAfterHeader > 0
- ? Math.min(retryAfterHeader, maxTimeout)
- : Math.min(minTimeout * timeoutFactor ** (counter - 1), maxTimeout)
- setTimeout(() => cb(null), retryTimeout)
- }
- onResponseStart (controller, statusCode, headers, statusMessage) {
- this.error = null
- this.retryCount += 1
- if (statusCode >= 300) {
- const err = new RequestRetryError('Request failed', statusCode, {
- headers,
- data: {
- count: this.retryCount
- }
- })
- this.onResponseStartWithRetry(controller, statusCode, headers, statusMessage, err)
- return
- }
- // Checkpoint for resume from where we left it
- if (this.headersSent) {
- // Only Partial Content 206 supposed to provide Content-Range,
- // any other status code that partially consumed the payload
- // should not be retried because it would result in downstream
- // wrongly concatenate multiple responses.
- if (statusCode !== 206 && (this.start > 0 || statusCode !== 200)) {
- throw new RequestRetryError('server does not support the range header and the payload was partially consumed', statusCode, {
- headers,
- data: { count: this.retryCount }
- })
- }
- const contentRange = parseRangeHeader(headers['content-range'])
- // If no content range
- if (!contentRange) {
- // We always throw here as we want to indicate that we entred unexpected path
- throw new RequestRetryError('Content-Range mismatch', statusCode, {
- headers,
- data: { count: this.retryCount }
- })
- }
- // Let's start with a weak etag check
- if (this.etag != null && this.etag !== headers.etag) {
- // We always throw here as we want to indicate that we entred unexpected path
- throw new RequestRetryError('ETag mismatch', statusCode, {
- headers,
- data: { count: this.retryCount }
- })
- }
- validatePartialResponseContentLength(headers, contentRange, statusCode, this.retryCount)
- const { start, size, end = size ? size - 1 : null } = contentRange
- if (this.start !== start || (this.end != null && this.end !== end)) {
- throw new RequestRetryError('Content-Range mismatch', statusCode, {
- headers,
- data: { count: this.retryCount }
- })
- }
- return
- }
- if (this.end == null) {
- if (statusCode === 206) {
- // First time we receive 206
- const range = parseRangeHeader(headers['content-range'])
- if (range == null) {
- this.headersSent = true
- this.handler.onResponseStart?.(
- this.controllerProxy,
- statusCode,
- headers,
- statusMessage
- )
- return
- }
- validatePartialResponseContentLength(headers, range, statusCode, this.retryCount)
- const { start, size, end = size ? size - 1 : null } = range
- assert(
- start != null && Number.isFinite(start),
- 'content-range mismatch'
- )
- assert(end != null && Number.isFinite(end), 'invalid content-length')
- this.start = start
- this.end = end
- }
- // We make our best to checkpoint the body for further range headers
- if (this.end == null) {
- const contentLength = headers['content-length']
- this.end = contentLength != null ? Number(contentLength) - 1 : null
- }
- assert(Number.isFinite(this.start))
- assert(
- this.end == null || Number.isFinite(this.end),
- 'invalid content-length'
- )
- this.resume = true
- this.etag = headers.etag != null ? headers.etag : null
- // Weak etags are not useful for comparison nor cache
- // for instance not safe to assume if the response is byte-per-byte
- // equal
- if (
- this.etag != null &&
- this.etag[0] === 'W' &&
- this.etag[1] === '/'
- ) {
- this.etag = null
- }
- this.headersSent = true
- this.handler.onResponseStart?.(
- this.controllerProxy,
- statusCode,
- headers,
- statusMessage
- )
- } else {
- throw new RequestRetryError('Request failed', statusCode, {
- headers,
- data: { count: this.retryCount }
- })
- }
- }
- onResponseData (_controller, chunk) {
- if (this.error) {
- return
- }
- this.start += chunk.length
- this.handler.onResponseData?.(this.controllerProxy, chunk)
- }
- onResponseEnd (_controller, trailers) {
- if (this.error && this.retryOpts.throwOnError) {
- throw this.error
- }
- if (!this.error) {
- this.retryCount = 0
- return this.handler.onResponseEnd?.(this.controllerProxy, trailers)
- }
- this.retry()
- }
- retry () {
- if (this.start !== 0) {
- const headers = { range: `bytes=${this.start}-${this.end ?? ''}` }
- // Weak etag check - weak etags will make comparison algorithms never match
- if (this.etag != null) {
- headers['if-match'] = this.etag
- }
- this.opts = {
- ...this.opts,
- headers: {
- ...this.opts.headers,
- ...headers
- }
- }
- }
- try {
- this.retryCountCheckpoint = this.retryCount
- this.dispatch(this.opts, this)
- } catch (err) {
- this.handler.onResponseError?.(this.controllerProxy, err)
- }
- }
- onResponseError (controller, err) {
- if (controller?.aborted || isDisturbed(this.opts.body) || (this.headersSent && !this.resume)) {
- this.handler.onResponseError?.(this.controllerProxy, err)
- return
- }
- function shouldRetry (returnedErr) {
- if (!returnedErr) {
- this.retry()
- return
- }
- this.handler?.onResponseError?.(this.controllerProxy, returnedErr)
- }
- // We reconcile in case of a mix between network errors
- // and server error response
- if (this.retryCount - this.retryCountCheckpoint > 0) {
- // We count the difference between the last checkpoint and the current retry count
- this.retryCount =
- this.retryCountCheckpoint +
- (this.retryCount - this.retryCountCheckpoint)
- } else {
- this.retryCount += 1
- }
- this.retryOpts.retry(
- err,
- {
- state: { counter: this.retryCount },
- opts: { retryOptions: this.retryOpts, ...this.opts }
- },
- shouldRetry.bind(this)
- )
- }
- }
- module.exports = RetryHandler
|