| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460 |
- 'use strict'
- const { RequestAbortedError } = require('../core/errors')
- /**
- * @typedef {import('../../types/dispatcher.d.ts').default.DispatchHandler} DispatchHandler
- */
- const DEFAULT_MAX_BUFFER_SIZE = 5 * 1024 * 1024
- /**
- * @typedef {Object} WaitingHandler
- * @property {DispatchHandler} handler
- * @property {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @property {Buffer[]} bufferedChunks
- * @property {number} bufferedBytes
- * @property {object | null} pendingTrailers
- * @property {boolean} done
- */
- /**
- * Handler that forwards response events to multiple waiting handlers.
- * Used for request deduplication.
- *
- * @implements {DispatchHandler}
- */
- class DeduplicationHandler {
- /**
- * @type {DispatchHandler}
- */
- #primaryHandler
- /**
- * @type {WaitingHandler[]}
- */
- #waitingHandlers = []
- /**
- * @type {number}
- */
- #maxBufferSize = DEFAULT_MAX_BUFFER_SIZE
- /**
- * @type {number}
- */
- #statusCode = 0
- /**
- * @type {Record<string, string | string[]>}
- */
- #headers = {}
- /**
- * @type {string}
- */
- #statusMessage = ''
- /**
- * @type {boolean}
- */
- #aborted = false
- /**
- * @type {boolean}
- */
- #responseStarted = false
- /**
- * @type {boolean}
- */
- #responseDataStarted = false
- /**
- * @type {boolean}
- */
- #completed = false
- /**
- * @type {import('../../types/dispatcher.d.ts').default.DispatchController | null}
- */
- #controller = null
- /**
- * @type {(() => void) | null}
- */
- #onComplete = null
- /**
- * @param {DispatchHandler} primaryHandler The primary handler
- * @param {() => void} onComplete Callback when request completes
- * @param {number} [maxBufferSize] Maximum paused buffer size per waiting handler
- */
- constructor (primaryHandler, onComplete, maxBufferSize = DEFAULT_MAX_BUFFER_SIZE) {
- this.#primaryHandler = primaryHandler
- this.#onComplete = onComplete
- this.#maxBufferSize = maxBufferSize
- }
- /**
- * Add a waiting handler that will receive response events.
- * Returns false if deduplication can no longer safely attach this handler.
- *
- * @param {DispatchHandler} handler
- * @returns {boolean}
- */
- addWaitingHandler (handler) {
- if (this.#completed || this.#responseDataStarted) {
- return false
- }
- const waitingHandler = this.#createWaitingHandler(handler)
- const waitingController = waitingHandler.controller
- try {
- handler.onRequestStart?.(waitingController, null)
- if (waitingController.aborted) {
- waitingHandler.done = true
- return true
- }
- if (this.#responseStarted) {
- handler.onResponseStart?.(
- waitingController,
- this.#statusCode,
- this.#headers,
- this.#statusMessage
- )
- }
- } catch {
- // Ignore errors from waiting handlers
- waitingHandler.done = true
- return true
- }
- if (!waitingController.aborted) {
- this.#waitingHandlers.push(waitingHandler)
- }
- return true
- }
- /**
- * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @param {any} context
- */
- onRequestStart (controller, context) {
- this.#controller = controller
- this.#primaryHandler.onRequestStart?.(controller, context)
- }
- /**
- * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @param {number} statusCode
- * @param {import('../../types/header.d.ts').IncomingHttpHeaders} headers
- * @param {Socket} socket
- */
- onRequestUpgrade (controller, statusCode, headers, socket) {
- this.#primaryHandler.onRequestUpgrade?.(controller, statusCode, headers, socket)
- }
- /**
- * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @param {number} statusCode
- * @param {Record<string, string | string[]>} headers
- * @param {string} statusMessage
- */
- onResponseStart (controller, statusCode, headers, statusMessage) {
- this.#responseStarted = true
- this.#statusCode = statusCode
- this.#headers = headers
- this.#statusMessage = statusMessage
- this.#primaryHandler.onResponseStart?.(controller, statusCode, headers, statusMessage)
- for (const waitingHandler of this.#waitingHandlers) {
- const { handler, controller: waitingController } = waitingHandler
- if (waitingHandler.done || waitingController.aborted) {
- waitingHandler.done = true
- continue
- }
- try {
- handler.onResponseStart?.(
- waitingController,
- statusCode,
- headers,
- statusMessage
- )
- } catch {
- // Ignore errors from waiting handlers
- }
- if (waitingController.aborted) {
- waitingHandler.done = true
- }
- }
- this.#pruneDoneWaitingHandlers()
- }
- /**
- * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @param {Buffer} chunk
- */
- onResponseData (controller, chunk) {
- if (this.#aborted || this.#completed) {
- return
- }
- this.#responseDataStarted = true
- this.#primaryHandler.onResponseData?.(controller, chunk)
- for (const waitingHandler of this.#waitingHandlers) {
- const { handler, controller: waitingController } = waitingHandler
- if (waitingHandler.done || waitingController.aborted) {
- waitingHandler.done = true
- continue
- }
- if (waitingController.paused) {
- this.#bufferWaitingChunk(waitingHandler, chunk)
- continue
- }
- try {
- handler.onResponseData?.(waitingController, chunk)
- } catch {
- // Ignore errors from waiting handlers
- }
- if (waitingController.aborted) {
- waitingHandler.done = true
- waitingHandler.bufferedChunks = []
- waitingHandler.bufferedBytes = 0
- }
- }
- this.#pruneDoneWaitingHandlers()
- }
- /**
- * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @param {object} trailers
- */
- onResponseEnd (controller, trailers) {
- if (this.#aborted || this.#completed) {
- return
- }
- this.#completed = true
- this.#primaryHandler.onResponseEnd?.(controller, trailers)
- for (const waitingHandler of this.#waitingHandlers) {
- if (waitingHandler.done || waitingHandler.controller.aborted) {
- waitingHandler.done = true
- continue
- }
- this.#flushWaitingHandler(waitingHandler)
- if (waitingHandler.done || waitingHandler.controller.aborted) {
- waitingHandler.done = true
- continue
- }
- if (waitingHandler.controller.paused && waitingHandler.bufferedChunks.length > 0) {
- waitingHandler.pendingTrailers = trailers
- continue
- }
- try {
- waitingHandler.handler.onResponseEnd?.(waitingHandler.controller, trailers)
- } catch {
- // Ignore errors from waiting handlers
- }
- waitingHandler.done = true
- }
- this.#pruneDoneWaitingHandlers()
- this.#onComplete?.()
- }
- /**
- * @param {import('../../types/dispatcher.d.ts').default.DispatchController} controller
- * @param {Error} err
- */
- onResponseError (controller, err) {
- if (this.#completed) {
- return
- }
- this.#aborted = true
- this.#completed = true
- this.#primaryHandler.onResponseError?.(controller, err)
- for (const waitingHandler of this.#waitingHandlers) {
- this.#errorWaitingHandler(waitingHandler, err)
- }
- this.#waitingHandlers = []
- this.#onComplete?.()
- }
- /**
- * @param {DispatchHandler} handler
- * @returns {WaitingHandler}
- */
- #createWaitingHandler (handler) {
- /** @type {WaitingHandler} */
- const waitingHandler = {
- handler,
- controller: null,
- bufferedChunks: [],
- bufferedBytes: 0,
- pendingTrailers: null,
- done: false
- }
- const state = {
- aborted: false,
- paused: false,
- reason: null
- }
- waitingHandler.controller = {
- resume: () => {
- if (state.aborted) {
- return
- }
- state.paused = false
- this.#flushWaitingHandler(waitingHandler)
- if (
- this.#completed &&
- waitingHandler.pendingTrailers &&
- waitingHandler.bufferedChunks.length === 0 &&
- !state.paused &&
- !state.aborted
- ) {
- try {
- waitingHandler.handler.onResponseEnd?.(waitingHandler.controller, waitingHandler.pendingTrailers)
- } catch {
- // Ignore errors from waiting handlers
- }
- waitingHandler.pendingTrailers = null
- waitingHandler.done = true
- }
- this.#pruneDoneWaitingHandlers()
- },
- pause: () => {
- if (!state.aborted) {
- state.paused = true
- }
- },
- get paused () { return state.paused },
- get aborted () { return state.aborted },
- get reason () { return state.reason },
- abort: (reason) => {
- state.aborted = true
- state.reason = reason ?? null
- waitingHandler.done = true
- waitingHandler.pendingTrailers = null
- waitingHandler.bufferedChunks = []
- waitingHandler.bufferedBytes = 0
- }
- }
- return waitingHandler
- }
- /**
- * @param {WaitingHandler} waitingHandler
- * @param {Buffer} chunk
- */
- #bufferWaitingChunk (waitingHandler, chunk) {
- if (waitingHandler.done || waitingHandler.controller.aborted) {
- waitingHandler.done = true
- waitingHandler.bufferedChunks = []
- waitingHandler.bufferedBytes = 0
- return
- }
- const bufferedChunk = Buffer.from(chunk)
- waitingHandler.bufferedChunks.push(bufferedChunk)
- waitingHandler.bufferedBytes += bufferedChunk.length
- if (waitingHandler.bufferedBytes > this.#maxBufferSize) {
- const err = new RequestAbortedError(`Deduplicated waiting handler exceeded maxBufferSize (${this.#maxBufferSize} bytes) while paused`)
- this.#errorWaitingHandler(waitingHandler, err)
- }
- }
- /**
- * @param {WaitingHandler} waitingHandler
- */
- #flushWaitingHandler (waitingHandler) {
- const { handler, controller } = waitingHandler
- while (
- !waitingHandler.done &&
- !controller.aborted &&
- !controller.paused &&
- waitingHandler.bufferedChunks.length > 0
- ) {
- const bufferedChunk = waitingHandler.bufferedChunks.shift()
- waitingHandler.bufferedBytes -= bufferedChunk.length
- try {
- handler.onResponseData?.(controller, bufferedChunk)
- } catch {
- // Ignore errors from waiting handlers
- }
- if (controller.aborted) {
- waitingHandler.done = true
- waitingHandler.pendingTrailers = null
- waitingHandler.bufferedChunks = []
- waitingHandler.bufferedBytes = 0
- break
- }
- }
- }
- /**
- * @param {WaitingHandler} waitingHandler
- * @param {Error} err
- */
- #errorWaitingHandler (waitingHandler, err) {
- if (waitingHandler.done) {
- return
- }
- waitingHandler.done = true
- waitingHandler.pendingTrailers = null
- waitingHandler.bufferedChunks = []
- waitingHandler.bufferedBytes = 0
- try {
- waitingHandler.controller.abort(err)
- waitingHandler.handler.onResponseError?.(waitingHandler.controller, err)
- } catch {
- // Ignore errors from waiting handlers
- }
- }
- #pruneDoneWaitingHandlers () {
- this.#waitingHandlers = this.#waitingHandlers.filter(waitingHandler => waitingHandler.done === false)
- }
- }
- module.exports = DeduplicationHandler
|