decompress.js 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551
  1. 'use strict'
  2. const { createInflate, createGunzip, createBrotliDecompress, createZstdDecompress } = require('node:zlib')
  3. const { pipeline, Transform: TransformStream } = require('node:stream')
  4. const { InvalidArgumentError, ResponseExceededMaxSizeError } = require('../core/errors')
  5. const DecoratorHandler = require('../handler/decorator-handler')
  6. const { runtimeFeatures } = require('../util/runtime-features')
  7. /** @typedef {import('node:stream').Transform} Transform */
  8. /** @typedef {import('node:stream').Transform} Controller */
  9. /** @typedef {Transform&import('node:zlib').Zlib} DecompressorStream */
  10. class DecompressController {
  11. #onPause
  12. #onResume
  13. #onAbort
  14. #paused = false
  15. constructor (onPause, onResume, onAbort) {
  16. this.#onPause = onPause
  17. this.#onResume = onResume
  18. this.#onAbort = onAbort
  19. this.target = null
  20. }
  21. pause () {
  22. if (this.#paused) {
  23. return
  24. }
  25. this.#paused = true
  26. this.#onPause()
  27. }
  28. resume () {
  29. if (!this.#paused) {
  30. return
  31. }
  32. this.#paused = false
  33. this.#onResume()
  34. }
  35. abort (reason) {
  36. this.target?.abort(reason)
  37. this.#onAbort(reason)
  38. }
  39. get paused () { return this.#paused }
  40. get aborted () { return this.target?.aborted ?? false }
  41. get reason () { return this.target?.reason ?? null }
  42. get rawHeaders () { return this.target?.rawHeaders ?? null }
  43. set rawHeaders (value) {
  44. if (this.target) {
  45. this.target.rawHeaders = value
  46. }
  47. }
  48. get rawTrailers () { return this.target?.rawTrailers ?? null }
  49. set rawTrailers (value) {
  50. if (this.target) {
  51. this.target.rawTrailers = value
  52. }
  53. }
  54. }
  55. /** @type {Record<string, () => DecompressorStream>} */
  56. const supportedEncodings = {
  57. gzip: createGunzip,
  58. 'x-gzip': createGunzip,
  59. br: createBrotliDecompress,
  60. deflate: createInflate,
  61. compress: createInflate,
  62. 'x-compress': createInflate,
  63. ...(runtimeFeatures.has('zstd') ? { zstd: createZstdDecompress } : {})
  64. }
  65. const defaultSkipStatusCodes = /** @type {const} */ ([204, 304])
  66. const defaultMaxSize = 0
  67. /**
  68. * Limits the output of one stage in a decompression chain.
  69. * @param {number} maxSize - Maximum output size in bytes
  70. * @returns {Transform}
  71. */
  72. function createMaxSizeLimiter (maxSize) {
  73. let size = 0
  74. return new TransformStream({
  75. transform (chunk, _encoding, callback) {
  76. const decompressedSize = size + chunk.length
  77. if (decompressedSize > maxSize) {
  78. callback(new ResponseExceededMaxSizeError(
  79. `Decompressed response size (${decompressedSize}) exceeded maxSize (${maxSize})`
  80. ))
  81. return
  82. }
  83. size = decompressedSize
  84. callback(null, chunk)
  85. }
  86. })
  87. }
  88. let warningEmitted = /** @type {boolean} */ (false)
  89. /**
  90. * @typedef {Object} DecompressHandlerOptions
  91. * @property {number[]|Readonly<number[]>} [skipStatusCodes=[204, 304]] - List of status codes to skip decompression for
  92. * @property {boolean} [skipErrorResponses] - Whether to skip decompression for error responses (status codes >= 400)
  93. * @property {number} [maxSize=0] - Maximum decompressed response size in bytes. 0 disables the limit
  94. */
  95. class DecompressHandler extends DecoratorHandler {
  96. /** @type {Transform[]} */
  97. #decompressors = []
  98. /** @type {Record<string, string | string[]> | undefined} */
  99. #trailers
  100. /** @type {Readonly<number[]>} */
  101. #skipStatusCodes
  102. /** @type {boolean} */
  103. #skipErrorResponses
  104. /** @type {number} */
  105. #maxSize
  106. /** @type {number} */
  107. #decompressedSize = 0
  108. /** @type {boolean} */
  109. #terminated = false
  110. /** @type {boolean} */
  111. #inputEnded = false
  112. /** @type {boolean} */
  113. #inputBackpressured = false
  114. /** @type {boolean} */
  115. #upstreamPaused = false
  116. /** @type {boolean} */
  117. #draining = false
  118. /** @type {boolean} */
  119. #drainRequested = false
  120. /** @type {boolean} */
  121. #completionPending = false
  122. /** @type {DecompressorStream | undefined} */
  123. #finalDecompressor
  124. /** @type {DecompressController} */
  125. #controller
  126. constructor (handler, { skipStatusCodes = defaultSkipStatusCodes, skipErrorResponses = true, maxSize = defaultMaxSize } = {}) {
  127. if (!Number.isSafeInteger(maxSize) || maxSize < 0) {
  128. throw new InvalidArgumentError('maxSize must be a non-negative integer')
  129. }
  130. super(handler)
  131. this.#skipStatusCodes = skipStatusCodes
  132. this.#skipErrorResponses = skipErrorResponses
  133. this.#maxSize = maxSize
  134. this.#controller = new DecompressController(
  135. () => this.#onDownstreamPause(),
  136. () => this.#onDownstreamResume(),
  137. reason => {
  138. if (this.#inputEnded && !this.#terminated) {
  139. this.onResponseError(this.#controller, reason)
  140. }
  141. }
  142. )
  143. }
  144. #onDownstreamPause () {
  145. this.#pauseUpstream()
  146. }
  147. #onDownstreamResume () {
  148. const drainWasDeferred = this.#draining
  149. this.#drainOutput()
  150. if (!drainWasDeferred) {
  151. this.#resumeUpstreamIfNeeded()
  152. this.#finishIfReady()
  153. }
  154. }
  155. #pauseUpstream () {
  156. if (!this.#upstreamPaused && !this.#terminated) {
  157. this.#upstreamPaused = true
  158. this.#controller.target?.pause()
  159. }
  160. }
  161. #resumeUpstreamIfNeeded () {
  162. if (this.#upstreamPaused && !this.#controller.paused && !this.#inputBackpressured) {
  163. this.#upstreamPaused = false
  164. if (!this.#inputEnded) {
  165. this.#controller.target?.resume()
  166. }
  167. }
  168. }
  169. #drainOutput () {
  170. if (this.#terminated || this.#controller.paused || !this.#finalDecompressor) {
  171. return
  172. }
  173. if (this.#draining) {
  174. this.#drainRequested = true
  175. return
  176. }
  177. this.#draining = true
  178. try {
  179. do {
  180. this.#drainRequested = false
  181. let chunk
  182. while (!this.#terminated && !this.#controller.paused && (chunk = this.#finalDecompressor.read()) !== null) {
  183. if (this.#maxSize > 0) {
  184. const decompressedSize = this.#decompressedSize + chunk.length
  185. if (decompressedSize > this.#maxSize) {
  186. this.#fail(new ResponseExceededMaxSizeError(
  187. `Decompressed response size (${decompressedSize}) exceeded maxSize (${this.#maxSize})`
  188. ))
  189. return
  190. }
  191. this.#decompressedSize = decompressedSize
  192. }
  193. const result = super.onResponseData(this.#controller, chunk)
  194. if (result === false && !this.#controller.paused) {
  195. this.#controller.pause()
  196. }
  197. }
  198. } while (this.#drainRequested && !this.#terminated && !this.#controller.paused)
  199. } finally {
  200. this.#draining = false
  201. }
  202. this.#resumeUpstreamIfNeeded()
  203. this.#finishIfReady()
  204. }
  205. #finishIfReady () {
  206. if (this.#terminated || !this.#completionPending || this.#controller.paused || this.#draining) {
  207. return
  208. }
  209. this.#terminated = true
  210. this.#cleanupDecompressors()
  211. super.onResponseEnd(this.#controller, this.#trailers)
  212. }
  213. #onDecompressionEnd () {
  214. if (this.#terminated) {
  215. return
  216. }
  217. this.#completionPending = true
  218. this.#drainOutput()
  219. this.#finishIfReady()
  220. }
  221. /**
  222. * Determines if decompression should be skipped based on encoding and status code
  223. * @param {string} contentEncoding - Content-Encoding header value
  224. * @param {number} statusCode - HTTP status code of the response
  225. * @returns {boolean} - True if decompression should be skipped
  226. */
  227. #shouldSkipDecompression (contentEncoding, statusCode) {
  228. if (!contentEncoding || statusCode < 200) return true
  229. if (this.#skipStatusCodes.includes(statusCode)) return true
  230. if (this.#skipErrorResponses && statusCode >= 400) return true
  231. return false
  232. }
  233. /**
  234. * Creates a chain of decompressors for multiple content encodings
  235. *
  236. * @param {string} encodings - Comma-separated list of content encodings
  237. * @returns {Array<Transform>} - Array of decompressor and limiting streams
  238. * @throws {Error} - If the number of content-encodings exceeds the maximum allowed
  239. */
  240. #createDecompressionChain (encodings) {
  241. const parts = encodings.split(',')
  242. // Limit the number of content-encodings to prevent resource exhaustion.
  243. // CVE fix similar to urllib3 (GHSA-gm62-xv2j-4w53) and curl (CVE-2022-32206).
  244. const maxContentEncodings = 5
  245. if (parts.length > maxContentEncodings) {
  246. throw new Error(`too many content-encodings in response: ${parts.length}, maximum allowed is ${maxContentEncodings}`)
  247. }
  248. /** @type {DecompressorStream[]} */
  249. const decompressors = []
  250. for (let i = parts.length - 1; i >= 0; i--) {
  251. const encoding = parts[i].trim()
  252. if (!encoding) continue
  253. if (!supportedEncodings[encoding]) {
  254. decompressors.length = 0 // Clear if unsupported encoding
  255. return decompressors // Unsupported encoding
  256. }
  257. decompressors.push(supportedEncodings[encoding]())
  258. }
  259. if (decompressors.length < 2) {
  260. return decompressors
  261. }
  262. /** @type {Transform[]} */
  263. const streams = []
  264. for (let i = 0; i < decompressors.length; i++) {
  265. streams.push(decompressors[i])
  266. if (i < decompressors.length - 1 && this.#maxSize > 0) {
  267. streams.push(createMaxSizeLimiter(this.#maxSize))
  268. }
  269. }
  270. return streams
  271. }
  272. /**
  273. * Stops decompression and reports an error.
  274. * @param {Error} error - The decompression error
  275. * @returns {void}
  276. */
  277. #fail (error) {
  278. if (this.#terminated) {
  279. return
  280. }
  281. if (this.#inputEnded) {
  282. // The request is already marked complete once the compressed input ends,
  283. // so controller.abort() can no longer propagate decoder flush errors.
  284. this.onResponseError(this.#controller, error)
  285. } else {
  286. this.#controller.abort(error)
  287. }
  288. }
  289. /**
  290. * Sets up event handlers for the final decompressor stream.
  291. * @param {DecompressorStream} decompressor - The decompressor stream
  292. * @returns {void}
  293. */
  294. #setupDecompressorEvents (decompressor) {
  295. this.#finalDecompressor = decompressor
  296. decompressor.on('readable', () => this.#drainOutput())
  297. decompressor.on('error', (error) => this.#fail(error))
  298. }
  299. /**
  300. * Sets up event handling for a single decompressor
  301. * @returns {void}
  302. */
  303. #setupSingleDecompressor () {
  304. const decompressor = this.#decompressors[0]
  305. this.#setupDecompressorEvents(decompressor)
  306. decompressor.on('end', () => this.#onDecompressionEnd())
  307. }
  308. /**
  309. * Sets up event handling for multiple chained decompressors using pipeline
  310. * @returns {void}
  311. */
  312. #setupMultipleDecompressors () {
  313. const lastDecompressor = this.#decompressors[this.#decompressors.length - 1]
  314. this.#setupDecompressorEvents(lastDecompressor)
  315. pipeline(this.#decompressors, (err) => {
  316. if (this.#terminated) {
  317. return
  318. }
  319. if (err) {
  320. this.#fail(err)
  321. return
  322. }
  323. this.#onDecompressionEnd()
  324. })
  325. }
  326. #setupInputBackpressure () {
  327. const decompressor = this.#decompressors[0]
  328. decompressor.on('drain', () => {
  329. if (this.#terminated) {
  330. return
  331. }
  332. this.#inputBackpressured = false
  333. if (!this.#controller.paused) {
  334. this.#drainOutput()
  335. this.#resumeUpstreamIfNeeded()
  336. }
  337. })
  338. }
  339. /**
  340. * Cleans up decompressor references to prevent memory leaks
  341. * @returns {void}
  342. */
  343. #cleanupDecompressors () {
  344. this.#decompressors.length = 0
  345. this.#finalDecompressor = undefined
  346. }
  347. onRequestStart (controller, context) {
  348. this.#controller.target = controller
  349. return super.onRequestStart(this.#controller, context)
  350. }
  351. onRequestUpgrade (controller, statusCode, headers, socket) {
  352. return super.onRequestUpgrade(this.#controller, statusCode, headers, socket)
  353. }
  354. /**
  355. * @param {Controller} controller
  356. * @param {number} statusCode
  357. * @param {Record<string, string | string[] | undefined>} headers
  358. * @param {string} statusMessage
  359. * @returns {void}
  360. */
  361. onResponseStart (controller, statusCode, headers, statusMessage) {
  362. const contentEncoding = headers['content-encoding']
  363. // If content encoding is not supported or status code is in skip list
  364. if (this.#shouldSkipDecompression(contentEncoding, statusCode)) {
  365. return super.onResponseStart(this.#controller, statusCode, headers, statusMessage)
  366. }
  367. const decompressors = this.#createDecompressionChain(contentEncoding.toLowerCase())
  368. if (decompressors.length === 0) {
  369. this.#cleanupDecompressors()
  370. return super.onResponseStart(this.#controller, statusCode, headers, statusMessage)
  371. }
  372. this.#decompressors = decompressors
  373. // Remove compression headers since we're decompressing
  374. const { 'content-encoding': _, 'content-length': __, ...newHeaders } = headers
  375. if (this.#controller.rawHeaders) {
  376. const rawHeaders = this.#controller.rawHeaders
  377. if (Array.isArray(rawHeaders)) {
  378. const filteredHeaders = []
  379. for (let i = 0; i < rawHeaders.length; i += 2) {
  380. const headerName = rawHeaders[i]
  381. const name = Buffer.isBuffer(headerName) ? headerName.toString('latin1') : `${headerName}`
  382. const lowerName = name.toLowerCase()
  383. if (lowerName === 'content-encoding' || lowerName === 'content-length') {
  384. continue
  385. }
  386. filteredHeaders.push(rawHeaders[i], rawHeaders[i + 1])
  387. }
  388. rawHeaders.splice(0, rawHeaders.length, ...filteredHeaders)
  389. } else if (typeof rawHeaders === 'object') {
  390. for (const name of Object.keys(rawHeaders)) {
  391. const lowerName = name.toLowerCase()
  392. if (lowerName === 'content-encoding' || lowerName === 'content-length') {
  393. delete rawHeaders[name]
  394. }
  395. }
  396. }
  397. }
  398. this.#setupInputBackpressure()
  399. if (this.#decompressors.length === 1) {
  400. this.#setupSingleDecompressor()
  401. } else {
  402. this.#setupMultipleDecompressors()
  403. }
  404. return super.onResponseStart(this.#controller, statusCode, newHeaders, statusMessage)
  405. }
  406. /**
  407. * @param {Controller} controller
  408. * @param {Buffer} chunk
  409. * @returns {void}
  410. */
  411. onResponseData (controller, chunk) {
  412. if (this.#decompressors.length > 0) {
  413. if (!this.#decompressors[0].write(chunk)) {
  414. this.#inputBackpressured = true
  415. this.#pauseUpstream()
  416. }
  417. return
  418. }
  419. return super.onResponseData(this.#controller, chunk)
  420. }
  421. /**
  422. * @param {Controller} controller
  423. * @param {Record<string, string | string[]> | undefined} trailers
  424. * @returns {void}
  425. */
  426. onResponseEnd (controller, trailers) {
  427. if (this.#decompressors.length > 0) {
  428. this.#inputEnded = true
  429. this.#trailers = trailers
  430. this.#decompressors[0].end()
  431. return
  432. }
  433. return super.onResponseEnd(this.#controller, trailers)
  434. }
  435. /**
  436. * @param {Controller} controller
  437. * @param {Error} err
  438. * @returns {void}
  439. */
  440. onResponseError (controller, err) {
  441. if (this.#terminated) {
  442. return
  443. }
  444. this.#terminated = true
  445. for (const decompressor of this.#decompressors) {
  446. decompressor.destroy()
  447. }
  448. this.#cleanupDecompressors()
  449. super.onResponseError(this.#controller, err)
  450. }
  451. }
  452. /**
  453. * Creates a decompression interceptor for HTTP responses
  454. * @param {DecompressHandlerOptions} [options] - Options for the interceptor
  455. * @returns {Function} - Interceptor function
  456. */
  457. function createDecompressInterceptor (options = {}) {
  458. // Emit experimental warning only once
  459. if (!warningEmitted) {
  460. process.emitWarning(
  461. 'DecompressInterceptor is experimental and subject to change',
  462. 'ExperimentalWarning'
  463. )
  464. warningEmitted = true
  465. }
  466. return (dispatch) => {
  467. return (opts, handler) => {
  468. const decompressHandler = new DecompressHandler(handler, options)
  469. return dispatch(opts, decompressHandler)
  470. }
  471. }
  472. }
  473. module.exports = createDecompressInterceptor