deduplication-handler.js 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460
  1. 'use strict'
  2. const { RequestAbortedError } = require('../core/errors')
  3. /**
  4. * @typedef {import('../../types/dispatcher.d.ts').default.DispatchHandler} DispatchHandler
  5. */
  6. const DEFAULT_MAX_BUFFER_SIZE = 5 * 1024 * 1024
  7. /**
  8. * @typedef {Object} WaitingHandler
  9. * @property {DispatchHandler} handler
  10. * @property {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  11. * @property {Buffer[]} bufferedChunks
  12. * @property {number} bufferedBytes
  13. * @property {object | null} pendingTrailers
  14. * @property {boolean} done
  15. */
  16. /**
  17. * Handler that forwards response events to multiple waiting handlers.
  18. * Used for request deduplication.
  19. *
  20. * @implements {DispatchHandler}
  21. */
  22. class DeduplicationHandler {
  23. /**
  24. * @type {DispatchHandler}
  25. */
  26. #primaryHandler
  27. /**
  28. * @type {WaitingHandler[]}
  29. */
  30. #waitingHandlers = []
  31. /**
  32. * @type {number}
  33. */
  34. #maxBufferSize = DEFAULT_MAX_BUFFER_SIZE
  35. /**
  36. * @type {number}
  37. */
  38. #statusCode = 0
  39. /**
  40. * @type {Record<string, string | string[]>}
  41. */
  42. #headers = {}
  43. /**
  44. * @type {string}
  45. */
  46. #statusMessage = ''
  47. /**
  48. * @type {boolean}
  49. */
  50. #aborted = false
  51. /**
  52. * @type {boolean}
  53. */
  54. #responseStarted = false
  55. /**
  56. * @type {boolean}
  57. */
  58. #responseDataStarted = false
  59. /**
  60. * @type {boolean}
  61. */
  62. #completed = false
  63. /**
  64. * @type {import('../../types/dispatcher.d.ts').default.DispatchController | null}
  65. */
  66. #controller = null
  67. /**
  68. * @type {(() => void) | null}
  69. */
  70. #onComplete = null
  71. /**
  72. * @param {DispatchHandler} primaryHandler The primary handler
  73. * @param {() => void} onComplete Callback when request completes
  74. * @param {number} [maxBufferSize] Maximum paused buffer size per waiting handler
  75. */
  76. constructor (primaryHandler, onComplete, maxBufferSize = DEFAULT_MAX_BUFFER_SIZE) {
  77. this.#primaryHandler = primaryHandler
  78. this.#onComplete = onComplete
  79. this.#maxBufferSize = maxBufferSize
  80. }
  81. /**
  82. * Add a waiting handler that will receive response events.
  83. * Returns false if deduplication can no longer safely attach this handler.
  84. *
  85. * @param {DispatchHandler} handler
  86. * @returns {boolean}
  87. */
  88. addWaitingHandler (handler) {
  89. if (this.#completed || this.#responseDataStarted) {
  90. return false
  91. }
  92. const waitingHandler = this.#createWaitingHandler(handler)
  93. const waitingController = waitingHandler.controller
  94. try {
  95. handler.onRequestStart?.(waitingController, null)
  96. if (waitingController.aborted) {
  97. waitingHandler.done = true
  98. return true
  99. }
  100. if (this.#responseStarted) {
  101. handler.onResponseStart?.(
  102. waitingController,
  103. this.#statusCode,
  104. this.#headers,
  105. this.#statusMessage
  106. )
  107. }
  108. } catch {
  109. // Ignore errors from waiting handlers
  110. waitingHandler.done = true
  111. return true
  112. }
  113. if (!waitingController.aborted) {
  114. this.#waitingHandlers.push(waitingHandler)
  115. }
  116. return true
  117. }
  118. /**
  119. * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  120. * @param {any} context
  121. */
  122. onRequestStart (controller, context) {
  123. this.#controller = controller
  124. this.#primaryHandler.onRequestStart?.(controller, context)
  125. }
  126. /**
  127. * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  128. * @param {number} statusCode
  129. * @param {import('../../types/header.d.ts').IncomingHttpHeaders} headers
  130. * @param {Socket} socket
  131. */
  132. onRequestUpgrade (controller, statusCode, headers, socket) {
  133. this.#primaryHandler.onRequestUpgrade?.(controller, statusCode, headers, socket)
  134. }
  135. /**
  136. * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  137. * @param {number} statusCode
  138. * @param {Record<string, string | string[]>} headers
  139. * @param {string} statusMessage
  140. */
  141. onResponseStart (controller, statusCode, headers, statusMessage) {
  142. this.#responseStarted = true
  143. this.#statusCode = statusCode
  144. this.#headers = headers
  145. this.#statusMessage = statusMessage
  146. this.#primaryHandler.onResponseStart?.(controller, statusCode, headers, statusMessage)
  147. for (const waitingHandler of this.#waitingHandlers) {
  148. const { handler, controller: waitingController } = waitingHandler
  149. if (waitingHandler.done || waitingController.aborted) {
  150. waitingHandler.done = true
  151. continue
  152. }
  153. try {
  154. handler.onResponseStart?.(
  155. waitingController,
  156. statusCode,
  157. headers,
  158. statusMessage
  159. )
  160. } catch {
  161. // Ignore errors from waiting handlers
  162. }
  163. if (waitingController.aborted) {
  164. waitingHandler.done = true
  165. }
  166. }
  167. this.#pruneDoneWaitingHandlers()
  168. }
  169. /**
  170. * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  171. * @param {Buffer} chunk
  172. */
  173. onResponseData (controller, chunk) {
  174. if (this.#aborted || this.#completed) {
  175. return
  176. }
  177. this.#responseDataStarted = true
  178. this.#primaryHandler.onResponseData?.(controller, chunk)
  179. for (const waitingHandler of this.#waitingHandlers) {
  180. const { handler, controller: waitingController } = waitingHandler
  181. if (waitingHandler.done || waitingController.aborted) {
  182. waitingHandler.done = true
  183. continue
  184. }
  185. if (waitingController.paused) {
  186. this.#bufferWaitingChunk(waitingHandler, chunk)
  187. continue
  188. }
  189. try {
  190. handler.onResponseData?.(waitingController, chunk)
  191. } catch {
  192. // Ignore errors from waiting handlers
  193. }
  194. if (waitingController.aborted) {
  195. waitingHandler.done = true
  196. waitingHandler.bufferedChunks = []
  197. waitingHandler.bufferedBytes = 0
  198. }
  199. }
  200. this.#pruneDoneWaitingHandlers()
  201. }
  202. /**
  203. * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  204. * @param {object} trailers
  205. */
  206. onResponseEnd (controller, trailers) {
  207. if (this.#aborted || this.#completed) {
  208. return
  209. }
  210. this.#completed = true
  211. this.#primaryHandler.onResponseEnd?.(controller, trailers)
  212. for (const waitingHandler of this.#waitingHandlers) {
  213. if (waitingHandler.done || waitingHandler.controller.aborted) {
  214. waitingHandler.done = true
  215. continue
  216. }
  217. this.#flushWaitingHandler(waitingHandler)
  218. if (waitingHandler.done || waitingHandler.controller.aborted) {
  219. waitingHandler.done = true
  220. continue
  221. }
  222. if (waitingHandler.controller.paused && waitingHandler.bufferedChunks.length > 0) {
  223. waitingHandler.pendingTrailers = trailers
  224. continue
  225. }
  226. try {
  227. waitingHandler.handler.onResponseEnd?.(waitingHandler.controller, trailers)
  228. } catch {
  229. // Ignore errors from waiting handlers
  230. }
  231. waitingHandler.done = true
  232. }
  233. this.#pruneDoneWaitingHandlers()
  234. this.#onComplete?.()
  235. }
  236. /**
  237. * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
  238. * @param {Error} err
  239. */
  240. onResponseError (controller, err) {
  241. if (this.#completed) {
  242. return
  243. }
  244. this.#aborted = true
  245. this.#completed = true
  246. this.#primaryHandler.onResponseError?.(controller, err)
  247. for (const waitingHandler of this.#waitingHandlers) {
  248. this.#errorWaitingHandler(waitingHandler, err)
  249. }
  250. this.#waitingHandlers = []
  251. this.#onComplete?.()
  252. }
  253. /**
  254. * @param {DispatchHandler} handler
  255. * @returns {WaitingHandler}
  256. */
  257. #createWaitingHandler (handler) {
  258. /** @type {WaitingHandler} */
  259. const waitingHandler = {
  260. handler,
  261. controller: null,
  262. bufferedChunks: [],
  263. bufferedBytes: 0,
  264. pendingTrailers: null,
  265. done: false
  266. }
  267. const state = {
  268. aborted: false,
  269. paused: false,
  270. reason: null
  271. }
  272. waitingHandler.controller = {
  273. resume: () => {
  274. if (state.aborted) {
  275. return
  276. }
  277. state.paused = false
  278. this.#flushWaitingHandler(waitingHandler)
  279. if (
  280. this.#completed &&
  281. waitingHandler.pendingTrailers &&
  282. waitingHandler.bufferedChunks.length === 0 &&
  283. !state.paused &&
  284. !state.aborted
  285. ) {
  286. try {
  287. waitingHandler.handler.onResponseEnd?.(waitingHandler.controller, waitingHandler.pendingTrailers)
  288. } catch {
  289. // Ignore errors from waiting handlers
  290. }
  291. waitingHandler.pendingTrailers = null
  292. waitingHandler.done = true
  293. }
  294. this.#pruneDoneWaitingHandlers()
  295. },
  296. pause: () => {
  297. if (!state.aborted) {
  298. state.paused = true
  299. }
  300. },
  301. get paused () { return state.paused },
  302. get aborted () { return state.aborted },
  303. get reason () { return state.reason },
  304. abort: (reason) => {
  305. state.aborted = true
  306. state.reason = reason ?? null
  307. waitingHandler.done = true
  308. waitingHandler.pendingTrailers = null
  309. waitingHandler.bufferedChunks = []
  310. waitingHandler.bufferedBytes = 0
  311. }
  312. }
  313. return waitingHandler
  314. }
  315. /**
  316. * @param {WaitingHandler} waitingHandler
  317. * @param {Buffer} chunk
  318. */
  319. #bufferWaitingChunk (waitingHandler, chunk) {
  320. if (waitingHandler.done || waitingHandler.controller.aborted) {
  321. waitingHandler.done = true
  322. waitingHandler.bufferedChunks = []
  323. waitingHandler.bufferedBytes = 0
  324. return
  325. }
  326. const bufferedChunk = Buffer.from(chunk)
  327. waitingHandler.bufferedChunks.push(bufferedChunk)
  328. waitingHandler.bufferedBytes += bufferedChunk.length
  329. if (waitingHandler.bufferedBytes > this.#maxBufferSize) {
  330. const err = new RequestAbortedError(`Deduplicated waiting handler exceeded maxBufferSize (${this.#maxBufferSize} bytes) while paused`)
  331. this.#errorWaitingHandler(waitingHandler, err)
  332. }
  333. }
  334. /**
  335. * @param {WaitingHandler} waitingHandler
  336. */
  337. #flushWaitingHandler (waitingHandler) {
  338. const { handler, controller } = waitingHandler
  339. while (
  340. !waitingHandler.done &&
  341. !controller.aborted &&
  342. !controller.paused &&
  343. waitingHandler.bufferedChunks.length > 0
  344. ) {
  345. const bufferedChunk = waitingHandler.bufferedChunks.shift()
  346. waitingHandler.bufferedBytes -= bufferedChunk.length
  347. try {
  348. handler.onResponseData?.(controller, bufferedChunk)
  349. } catch {
  350. // Ignore errors from waiting handlers
  351. }
  352. if (controller.aborted) {
  353. waitingHandler.done = true
  354. waitingHandler.pendingTrailers = null
  355. waitingHandler.bufferedChunks = []
  356. waitingHandler.bufferedBytes = 0
  357. break
  358. }
  359. }
  360. }
  361. /**
  362. * @param {WaitingHandler} waitingHandler
  363. * @param {Error} err
  364. */
  365. #errorWaitingHandler (waitingHandler, err) {
  366. if (waitingHandler.done) {
  367. return
  368. }
  369. waitingHandler.done = true
  370. waitingHandler.pendingTrailers = null
  371. waitingHandler.bufferedChunks = []
  372. waitingHandler.bufferedBytes = 0
  373. try {
  374. waitingHandler.controller.abort(err)
  375. waitingHandler.handler.onResponseError?.(waitingHandler.controller, err)
  376. } catch {
  377. // Ignore errors from waiting handlers
  378. }
  379. }
  380. #pruneDoneWaitingHandlers () {
  381. this.#waitingHandlers = this.#waitingHandlers.filter(waitingHandler => waitingHandler.done === false)
  382. }
  383. }
  384. module.exports = DeduplicationHandler