'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