retry-handler.js 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504
  1. 'use strict'
  2. const assert = require('node:assert')
  3. const { kRetryHandlerDefaultRetry } = require('../core/symbols')
  4. const { RequestRetryError } = require('../core/errors')
  5. const WrapHandler = require('./wrap-handler')
  6. const {
  7. isDisturbed,
  8. parseRangeHeader,
  9. wrapRequestBody
  10. } = require('../core/util')
  11. function calculateRetryAfterHeader (retryAfter) {
  12. const retryTime = new Date(retryAfter).getTime()
  13. return isNaN(retryTime) ? 0 : retryTime - Date.now()
  14. }
  15. function validatePartialResponseContentLength (headers, range, statusCode, retryCount) {
  16. const contentLength = headers['content-length']
  17. if (contentLength == null) {
  18. return
  19. }
  20. if (!Number.isFinite(range.start) || !Number.isFinite(range.end)) {
  21. return
  22. }
  23. const length = Number(contentLength)
  24. const expectedLength = range.end - range.start + 1
  25. if (!Number.isFinite(length) || length !== expectedLength) {
  26. throw new RequestRetryError('Content-Length mismatch', statusCode, {
  27. headers,
  28. data: { count: retryCount }
  29. })
  30. }
  31. }
  32. // A stable controller handed to the downstream handler for the lifetime of the
  33. // request. Each transparent retry/resume is a separate dispatch with its own
  34. // connection controller. The proxy always forwards to the active connection
  35. // while preserving a downstream pause across controller replacement.
  36. class RetryController {
  37. #paused = false
  38. #target = null
  39. set target (target) {
  40. this.#target = target
  41. if (this.#paused) {
  42. target?.pause()
  43. }
  44. }
  45. get target () { return this.#target }
  46. pause () {
  47. this.#paused = true
  48. this.#target?.pause()
  49. }
  50. resume () {
  51. this.#paused = false
  52. this.#target?.resume()
  53. }
  54. abort (reason) {
  55. this.#target?.abort(reason)
  56. }
  57. get paused () { return this.#paused || (this.#target?.paused ?? false) }
  58. get aborted () { return this.#target?.aborted ?? false }
  59. get reason () { return this.#target?.reason ?? null }
  60. get rawHeaders () { return this.#target?.rawHeaders ?? null }
  61. set rawHeaders (value) {
  62. if (this.#target) {
  63. this.#target.rawHeaders = value
  64. }
  65. }
  66. get rawTrailers () { return this.#target?.rawTrailers ?? null }
  67. set rawTrailers (value) {
  68. if (this.#target) {
  69. this.#target.rawTrailers = value
  70. }
  71. }
  72. }
  73. class RetryHandler {
  74. constructor (opts, { dispatch, handler }) {
  75. const { retryOptions, ...dispatchOpts } = opts
  76. const {
  77. // Retry scoped
  78. retry: retryFn,
  79. maxRetries,
  80. maxTimeout,
  81. minTimeout,
  82. timeoutFactor,
  83. // Response scoped
  84. methods,
  85. errorCodes,
  86. retryAfter,
  87. statusCodes,
  88. throwOnError
  89. } = retryOptions ?? {}
  90. this.error = null
  91. this.dispatch = dispatch
  92. this.handler = WrapHandler.wrap(handler)
  93. this.opts = { ...dispatchOpts, body: wrapRequestBody(opts.body) }
  94. this.retryOpts = {
  95. throwOnError: throwOnError ?? true,
  96. retry: retryFn ?? RetryHandler[kRetryHandlerDefaultRetry],
  97. retryAfter: retryAfter ?? true,
  98. maxTimeout: maxTimeout ?? 30 * 1000, // 30s,
  99. minTimeout: minTimeout ?? 500, // .5s
  100. timeoutFactor: timeoutFactor ?? 2,
  101. maxRetries: maxRetries ?? 5,
  102. // What errors we should retry
  103. methods: methods ?? ['GET', 'HEAD', 'OPTIONS', 'PUT', 'DELETE', 'TRACE'],
  104. // Indicates which errors to retry
  105. statusCodes: statusCodes ?? [500, 502, 503, 504, 429],
  106. // List of errors to retry
  107. errorCodes: errorCodes ?? [
  108. 'ECONNRESET',
  109. 'ECONNREFUSED',
  110. 'ENOTFOUND',
  111. 'ENETDOWN',
  112. 'ENETUNREACH',
  113. 'EHOSTDOWN',
  114. 'EHOSTUNREACH',
  115. 'EPIPE',
  116. 'UND_ERR_SOCKET'
  117. ]
  118. }
  119. this.retryCount = 0
  120. this.retryCountCheckpoint = 0
  121. this.headersSent = false
  122. this.start = 0
  123. this.end = null
  124. this.etag = null
  125. this.controllerProxy = new RetryController()
  126. }
  127. onResponseStartWithRetry (controller, statusCode, headers, statusMessage, err) {
  128. if (this.retryOpts.throwOnError) {
  129. // Preserve old behavior for status codes that are not eligible for retry
  130. if (this.retryOpts.statusCodes.includes(statusCode) === false) {
  131. if (this.headersSent) {
  132. // The downstream handler already received the response from an
  133. // earlier attempt. Forwarding this response would replace the
  134. // downstream body and leave the original body pending forever.
  135. this.handler.onResponseError?.(this.controllerProxy, err)
  136. } else {
  137. this.headersSent = true
  138. this.checkpointResponseEnd(headers)
  139. this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage)
  140. }
  141. } else {
  142. this.error = err
  143. }
  144. return
  145. }
  146. if (isDisturbed(this.opts.body)) {
  147. this.headersSent = true
  148. this.checkpointResponseEnd(headers)
  149. this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage)
  150. return
  151. }
  152. function shouldRetry (passedErr) {
  153. if (passedErr) {
  154. if (this.headersSent) {
  155. // The downstream handler already received the response from an
  156. // earlier attempt. Forwarding this response would replace the
  157. // downstream body and leave the original body pending forever.
  158. this.handler.onResponseError?.(this.controllerProxy, passedErr)
  159. } else {
  160. this.headersSent = true
  161. this.checkpointResponseEnd(headers)
  162. this.handler.onResponseStart?.(this.controllerProxy, statusCode, headers, statusMessage)
  163. }
  164. controller.resume()
  165. return
  166. }
  167. this.error = err
  168. controller.resume()
  169. }
  170. controller.pause()
  171. this.retryOpts.retry(
  172. err,
  173. {
  174. state: { counter: this.retryCount },
  175. opts: { retryOptions: this.retryOpts, ...this.opts }
  176. },
  177. shouldRetry.bind(this)
  178. )
  179. }
  180. checkpointResponseEnd (headers) {
  181. if (this.end == null && this.opts.method !== 'HEAD') {
  182. const contentLength = headers['content-length']
  183. this.end = contentLength != null ? Number(contentLength) - 1 : null
  184. assert(
  185. this.end == null || Number.isFinite(this.end),
  186. 'invalid content-length'
  187. )
  188. this.resume = this.end != null
  189. }
  190. }
  191. onRequestStart (controller, context) {
  192. this.controllerProxy.target = controller
  193. if (!this.headersSent) {
  194. this.handler.onRequestStart?.(this.controllerProxy, context)
  195. }
  196. }
  197. onRequestUpgrade (_controller, statusCode, headers, socket) {
  198. this.handler.onRequestUpgrade?.(this.controllerProxy, statusCode, headers, socket)
  199. }
  200. static [kRetryHandlerDefaultRetry] (err, { state, opts }, cb) {
  201. const { statusCode, code, headers } = err
  202. const { method, retryOptions } = opts
  203. const {
  204. maxRetries,
  205. minTimeout,
  206. maxTimeout,
  207. timeoutFactor,
  208. statusCodes,
  209. errorCodes,
  210. methods
  211. } = retryOptions
  212. const { counter } = state
  213. // Any code that is not a Undici's originated and allowed to retry
  214. if (code && code !== 'UND_ERR_REQ_RETRY' && !errorCodes.includes(code)) {
  215. cb(err)
  216. return
  217. }
  218. // If a set of method are provided and the current method is not in the list
  219. if (Array.isArray(methods) && !methods.includes(method)) {
  220. cb(err)
  221. return
  222. }
  223. // If a set of status code are provided and the current status code is not in the list
  224. if (
  225. statusCode != null &&
  226. Array.isArray(statusCodes) &&
  227. !statusCodes.includes(statusCode)
  228. ) {
  229. cb(err)
  230. return
  231. }
  232. // If we reached the max number of retries
  233. if (counter > maxRetries) {
  234. cb(err)
  235. return
  236. }
  237. let retryAfterHeader = headers?.['retry-after']
  238. if (retryAfterHeader) {
  239. retryAfterHeader = Number(retryAfterHeader)
  240. retryAfterHeader = Number.isNaN(retryAfterHeader)
  241. ? calculateRetryAfterHeader(headers['retry-after'])
  242. : retryAfterHeader * 1e3 // Retry-After is in seconds
  243. }
  244. const retryTimeout =
  245. retryAfterHeader > 0
  246. ? Math.min(retryAfterHeader, maxTimeout)
  247. : Math.min(minTimeout * timeoutFactor ** (counter - 1), maxTimeout)
  248. setTimeout(() => cb(null), retryTimeout)
  249. }
  250. onResponseStart (controller, statusCode, headers, statusMessage) {
  251. this.error = null
  252. this.retryCount += 1
  253. if (statusCode >= 300) {
  254. const err = new RequestRetryError('Request failed', statusCode, {
  255. headers,
  256. data: {
  257. count: this.retryCount
  258. }
  259. })
  260. this.onResponseStartWithRetry(controller, statusCode, headers, statusMessage, err)
  261. return
  262. }
  263. // Checkpoint for resume from where we left it
  264. if (this.headersSent) {
  265. // Only Partial Content 206 supposed to provide Content-Range,
  266. // any other status code that partially consumed the payload
  267. // should not be retried because it would result in downstream
  268. // wrongly concatenate multiple responses.
  269. if (statusCode !== 206 && (this.start > 0 || statusCode !== 200)) {
  270. throw new RequestRetryError('server does not support the range header and the payload was partially consumed', statusCode, {
  271. headers,
  272. data: { count: this.retryCount }
  273. })
  274. }
  275. const contentRange = parseRangeHeader(headers['content-range'])
  276. // If no content range
  277. if (!contentRange) {
  278. // We always throw here as we want to indicate that we entred unexpected path
  279. throw new RequestRetryError('Content-Range mismatch', statusCode, {
  280. headers,
  281. data: { count: this.retryCount }
  282. })
  283. }
  284. // Let's start with a weak etag check
  285. if (this.etag != null && this.etag !== headers.etag) {
  286. // We always throw here as we want to indicate that we entred unexpected path
  287. throw new RequestRetryError('ETag mismatch', statusCode, {
  288. headers,
  289. data: { count: this.retryCount }
  290. })
  291. }
  292. validatePartialResponseContentLength(headers, contentRange, statusCode, this.retryCount)
  293. const { start, size, end = size ? size - 1 : null } = contentRange
  294. if (this.start !== start || (this.end != null && this.end !== end)) {
  295. throw new RequestRetryError('Content-Range mismatch', statusCode, {
  296. headers,
  297. data: { count: this.retryCount }
  298. })
  299. }
  300. return
  301. }
  302. if (this.end == null) {
  303. if (statusCode === 206) {
  304. // First time we receive 206
  305. const range = parseRangeHeader(headers['content-range'])
  306. if (range == null) {
  307. this.headersSent = true
  308. this.handler.onResponseStart?.(
  309. this.controllerProxy,
  310. statusCode,
  311. headers,
  312. statusMessage
  313. )
  314. return
  315. }
  316. validatePartialResponseContentLength(headers, range, statusCode, this.retryCount)
  317. const { start, size, end = size ? size - 1 : null } = range
  318. assert(
  319. start != null && Number.isFinite(start),
  320. 'content-range mismatch'
  321. )
  322. assert(end != null && Number.isFinite(end), 'invalid content-length')
  323. this.start = start
  324. this.end = end
  325. }
  326. // We make our best to checkpoint the body for further range headers
  327. if (this.end == null) {
  328. const contentLength = headers['content-length']
  329. this.end = contentLength != null ? Number(contentLength) - 1 : null
  330. }
  331. assert(Number.isFinite(this.start))
  332. assert(
  333. this.end == null || Number.isFinite(this.end),
  334. 'invalid content-length'
  335. )
  336. this.resume = true
  337. this.etag = headers.etag != null ? headers.etag : null
  338. // Weak etags are not useful for comparison nor cache
  339. // for instance not safe to assume if the response is byte-per-byte
  340. // equal
  341. if (
  342. this.etag != null &&
  343. this.etag[0] === 'W' &&
  344. this.etag[1] === '/'
  345. ) {
  346. this.etag = null
  347. }
  348. this.headersSent = true
  349. this.handler.onResponseStart?.(
  350. this.controllerProxy,
  351. statusCode,
  352. headers,
  353. statusMessage
  354. )
  355. } else {
  356. throw new RequestRetryError('Request failed', statusCode, {
  357. headers,
  358. data: { count: this.retryCount }
  359. })
  360. }
  361. }
  362. onResponseData (_controller, chunk) {
  363. if (this.error) {
  364. return
  365. }
  366. this.start += chunk.length
  367. this.handler.onResponseData?.(this.controllerProxy, chunk)
  368. }
  369. onResponseEnd (_controller, trailers) {
  370. if (this.error && this.retryOpts.throwOnError) {
  371. throw this.error
  372. }
  373. if (!this.error) {
  374. this.retryCount = 0
  375. return this.handler.onResponseEnd?.(this.controllerProxy, trailers)
  376. }
  377. this.retry()
  378. }
  379. retry () {
  380. if (this.start !== 0) {
  381. const headers = { range: `bytes=${this.start}-${this.end ?? ''}` }
  382. // Weak etag check - weak etags will make comparison algorithms never match
  383. if (this.etag != null) {
  384. headers['if-match'] = this.etag
  385. }
  386. this.opts = {
  387. ...this.opts,
  388. headers: {
  389. ...this.opts.headers,
  390. ...headers
  391. }
  392. }
  393. }
  394. try {
  395. this.retryCountCheckpoint = this.retryCount
  396. this.dispatch(this.opts, this)
  397. } catch (err) {
  398. this.handler.onResponseError?.(this.controllerProxy, err)
  399. }
  400. }
  401. onResponseError (controller, err) {
  402. if (controller?.aborted || isDisturbed(this.opts.body) || (this.headersSent && !this.resume)) {
  403. this.handler.onResponseError?.(this.controllerProxy, err)
  404. return
  405. }
  406. function shouldRetry (returnedErr) {
  407. if (!returnedErr) {
  408. this.retry()
  409. return
  410. }
  411. this.handler?.onResponseError?.(this.controllerProxy, returnedErr)
  412. }
  413. // We reconcile in case of a mix between network errors
  414. // and server error response
  415. if (this.retryCount - this.retryCountCheckpoint > 0) {
  416. // We count the difference between the last checkpoint and the current retry count
  417. this.retryCount =
  418. this.retryCountCheckpoint +
  419. (this.retryCount - this.retryCountCheckpoint)
  420. } else {
  421. this.retryCount += 1
  422. }
  423. this.retryOpts.retry(
  424. err,
  425. {
  426. state: { counter: this.retryCount },
  427. opts: { retryOptions: this.retryOpts, ...this.opts }
  428. },
  429. shouldRetry.bind(this)
  430. )
  431. }
  432. }
  433. module.exports = RetryHandler