deduplicate.js 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117
  1. 'use strict'
  2. const diagnosticsChannel = require('node:diagnostics_channel')
  3. const util = require('../core/util')
  4. const DeduplicationHandler = require('../handler/deduplication-handler')
  5. const { normalizeHeaders, makeCacheKey, makeDeduplicationKey } = require('../util/cache.js')
  6. const pendingRequestsChannel = diagnosticsChannel.channel('undici:request:pending-requests')
  7. /**
  8. * @param {import('../../types/interceptors.d.ts').default.DeduplicateInterceptorOpts} [opts]
  9. * @returns {import('../../types/dispatcher.d.ts').default.DispatcherComposeInterceptor}
  10. */
  11. module.exports = (opts = {}) => {
  12. const {
  13. methods = ['GET'],
  14. skipHeaderNames = [],
  15. excludeHeaderNames = [],
  16. maxBufferSize = 5 * 1024 * 1024
  17. } = opts
  18. if (typeof opts !== 'object' || opts === null) {
  19. throw new TypeError(`expected type of opts to be an Object, got ${opts === null ? 'null' : typeof opts}`)
  20. }
  21. if (!Array.isArray(methods)) {
  22. throw new TypeError(`expected opts.methods to be an array, got ${typeof methods}`)
  23. }
  24. for (const method of methods) {
  25. if (!util.safeHTTPMethods.includes(method)) {
  26. throw new TypeError(`expected opts.methods to only contain safe HTTP methods, got ${method}`)
  27. }
  28. }
  29. if (!Array.isArray(skipHeaderNames)) {
  30. throw new TypeError(`expected opts.skipHeaderNames to be an array, got ${typeof skipHeaderNames}`)
  31. }
  32. if (!Array.isArray(excludeHeaderNames)) {
  33. throw new TypeError(`expected opts.excludeHeaderNames to be an array, got ${typeof excludeHeaderNames}`)
  34. }
  35. if (!Number.isFinite(maxBufferSize) || maxBufferSize <= 0) {
  36. throw new TypeError(`expected opts.maxBufferSize to be a positive finite number, got ${maxBufferSize}`)
  37. }
  38. // Convert to lowercase Set for case-insensitive header matching
  39. const skipHeaderNamesSet = new Set(skipHeaderNames.map(name => name.toLowerCase()))
  40. // Convert to lowercase Set for case-insensitive header exclusion from deduplication key
  41. const excludeHeaderNamesSet = new Set(excludeHeaderNames.map(name => name.toLowerCase()))
  42. /**
  43. * Map of pending requests for deduplication
  44. * @type {Map<string, DeduplicationHandler>}
  45. */
  46. const pendingRequests = new Map()
  47. return dispatch => {
  48. return (opts, handler) => {
  49. if (!opts.origin || methods.includes(opts.method) === false) {
  50. return dispatch(opts, handler)
  51. }
  52. opts = {
  53. ...opts,
  54. headers: normalizeHeaders(opts)
  55. }
  56. // Skip deduplication if request contains any of the specified headers
  57. if (skipHeaderNamesSet.size > 0) {
  58. for (const headerName of Object.keys(opts.headers)) {
  59. if (skipHeaderNamesSet.has(headerName.toLowerCase())) {
  60. return dispatch(opts, handler)
  61. }
  62. }
  63. }
  64. const cacheKey = makeCacheKey(opts)
  65. const dedupeKey = makeDeduplicationKey(cacheKey, excludeHeaderNamesSet)
  66. // Check if there's already a pending request for this key
  67. const pendingHandler = pendingRequests.get(dedupeKey)
  68. if (pendingHandler) {
  69. // Add this handler to the waiting list when safe.
  70. // If body streaming has already started, this request must be sent independently.
  71. if (pendingHandler.addWaitingHandler(handler)) {
  72. return true
  73. }
  74. return dispatch(opts, handler)
  75. }
  76. // Create a new deduplication handler
  77. const deduplicationHandler = new DeduplicationHandler(
  78. handler,
  79. () => {
  80. // Clean up when request completes
  81. pendingRequests.delete(dedupeKey)
  82. if (pendingRequestsChannel.hasSubscribers) {
  83. pendingRequestsChannel.publish({ size: pendingRequests.size, key: dedupeKey, type: 'removed' })
  84. }
  85. },
  86. maxBufferSize
  87. )
  88. // Register the pending request
  89. pendingRequests.set(dedupeKey, deduplicationHandler)
  90. if (pendingRequestsChannel.hasSubscribers) {
  91. pendingRequestsChannel.publish({ size: pendingRequests.size, key: dedupeKey, type: 'added' })
  92. }
  93. return dispatch(opts, deduplicationHandler)
  94. }
  95. }
  96. }