| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent f3ee312 commit 9962459
19 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -3,10 +3,10 @@ | |||
| 3 | 3 | const Readable = require('./readable') | |
| 4 | 4 | const { | |
| 5 | 5 | InvalidArgumentError, | |
| 6 | - RequestAbortedError, | ||
| 7 | - ResponseStatusCodeError | ||
| 6 | + RequestAbortedError | ||
| 8 | 7 | } = require('../core/errors') | |
| 9 | 8 | const util = require('../core/util') | |
| 9 | + const { getResolveErrorBodyCallback } = require('./util') | ||
| 10 | 10 | const { AsyncResource } = require('async_hooks') | |
| 11 | 11 | const { addSignal, removeSignal } = require('./abort-signal') | |
| 12 | 12 | ||
@@ -16,13 +16,17 @@ class RequestHandler extends AsyncResource { | |||
| 16 | 16 | throw new InvalidArgumentError('invalid opts') | |
| 17 | 17 | } | |
| 18 | 18 | ||
| 19 | - const { signal, method, opaque, body, onInfo, responseHeaders, throwOnError } = opts | ||
| 19 | + const { signal, method, opaque, body, onInfo, responseHeaders, throwOnError, highWaterMark } = opts | ||
| 20 | 20 | ||
| 21 | 21 | try { | |
| 22 | 22 | if (typeof callback !== 'function') { | |
| 23 | 23 | throw new InvalidArgumentError('invalid callback') | |
| 24 | 24 | } | |
| 25 | 25 | ||
| 26 | + if (highWaterMark && (typeof highWaterMark !== 'number' || highWaterMark < 0)) { | ||
| 27 | + throw new InvalidArgumentError('invalid highWaterMark') | ||
| 28 | + } | ||
| 29 | + | ||
| 26 | 30 | if (signal && typeof signal.on !== 'function' && typeof signal.addEventListener !== 'function') { | |
| 27 | 31 | throw new InvalidArgumentError('signal must be an EventEmitter or EventTarget') | |
| 28 | 32 | } | |
@@ -53,6 +57,7 @@ class RequestHandler extends AsyncResource { | |||
| 53 | 57 | this.context = null | |
| 54 | 58 | this.onInfo = onInfo || null | |
| 55 | 59 | this.throwOnError = throwOnError | |
| 60 | + this.highWaterMark = highWaterMark | ||
| 56 | 61 | ||
| 57 | 62 | if (util.isStream(body)) { | |
| 58 | 63 | body.on('error', (err) => { | |
@@ -73,40 +78,39 @@ class RequestHandler extends AsyncResource { | |||
| 73 | 78 | } | |
| 74 | 79 | ||
| 75 | 80 | onHeaders (statusCode, rawHeaders, resume, statusMessage) { | |
| 76 | - const { callback, opaque, abort, context } = this | ||
| 81 | + const { callback, opaque, abort, context, responseHeaders, highWaterMark } = this | ||
| 82 | + | ||
| 83 | + const headers = responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 77 | 84 | ||
| 78 | 85 | if (statusCode < 200) { | |
| 79 | 86 | if (this.onInfo) { | |
| 80 | - const headers = this.responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 81 | 87 | this.onInfo({ statusCode, headers }) | |
| 82 | 88 | } | |
| 83 | 89 | return | |
| 84 | 90 | } | |
| 85 | 91 | ||
| 86 | - const parsedHeaders = util.parseHeaders(rawHeaders) | ||
| 92 | + const parsedHeaders = responseHeaders === 'raw' ? util.parseHeaders(rawHeaders) : headers | ||
| 87 | 93 | const contentType = parsedHeaders['content-type'] | |
| 88 | - const body = new Readable(resume, abort, contentType) | ||
| 94 | + const body = new Readable({ resume, abort, contentType, highWaterMark }) | ||
| 89 | 95 | ||
| 90 | 96 | this.callback = null | |
| 91 | 97 | this.res = body | |
| 92 | - const headers = this.responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 93 | 98 | ||
| 94 | 99 | if (callback !== null) { | |
| 95 | 100 | if (this.throwOnError && statusCode >= 400) { | |
| 96 | 101 | this.runInAsyncScope(getResolveErrorBodyCallback, null, | |
| 97 | 102 | { callback, body, contentType, statusCode, statusMessage, headers } | |
| 98 | 103 | ) | |
| 99 | - return | ||
| 104 | + } else { | ||
| 105 | + this.runInAsyncScope(callback, null, null, { | ||
| 106 | + statusCode, | ||
| 107 | + headers, | ||
| 108 | + trailers: this.trailers, | ||
| 109 | + opaque, | ||
| 110 | + body, | ||
| 111 | + context | ||
| 112 | + }) | ||
| 100 | 113 | } | |
| 101 | - | ||
| 102 | - this.runInAsyncScope(callback, null, null, { | ||
| 103 | - statusCode, | ||
| 104 | - headers, | ||
| 105 | - trailers: this.trailers, | ||
| 106 | - opaque, | ||
| 107 | - body, | ||
| 108 | - context | ||
| 109 | - }) | ||
| 110 | 114 | } | |
| 111 | 115 | } | |
| 112 | 116 | ||
@@ -153,33 +157,6 @@ class RequestHandler extends AsyncResource { | |||
| 153 | 157 | } | |
| 154 | 158 | } | |
| 155 | 159 | ||
| 156 | - async function getResolveErrorBodyCallback ({ callback, body, contentType, statusCode, statusMessage, headers }) { | ||
| 157 | - if (statusCode === 204 || !contentType) { | ||
| 158 | - body.dump() | ||
| 159 | - process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers)) | ||
| 160 | - return | ||
| 161 | - } | ||
| 162 | - | ||
| 163 | - try { | ||
| 164 | - if (contentType.startsWith('application/json')) { | ||
| 165 | - const payload = await body.json() | ||
| 166 | - process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers, payload)) | ||
| 167 | - return | ||
| 168 | - } | ||
| 169 | - | ||
| 170 | - if (contentType.startsWith('text/')) { | ||
| 171 | - const payload = await body.text() | ||
| 172 | - process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers, payload)) | ||
| 173 | - return | ||
| 174 | - } | ||
| 175 | - } catch (err) { | ||
| 176 | - // Process in a fallback if error | ||
| 177 | - } | ||
| 178 | - | ||
| 179 | - body.dump() | ||
| 180 | - process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers)) | ||
| 181 | - } | ||
| 182 | - | ||
| 183 | 160 | function request (opts, callback) { | |
| 184 | 161 | if (callback === undefined) { | |
| 185 | 162 | return new Promise((resolve, reject) => { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -4,10 +4,10 @@ const { finished, PassThrough } = require('stream') | |||
| 4 | 4 | const { | |
| 5 | 5 | InvalidArgumentError, | |
| 6 | 6 | InvalidReturnValueError, | |
| 7 | - RequestAbortedError, | ||
| 8 | - ResponseStatusCodeError | ||
| 7 | + RequestAbortedError | ||
| 9 | 8 | } = require('../core/errors') | |
| 10 | 9 | const util = require('../core/util') | |
| 10 | + const { getResolveErrorBodyCallback } = require('./util') | ||
| 11 | 11 | const { AsyncResource } = require('async_hooks') | |
| 12 | 12 | const { addSignal, removeSignal } = require('./abort-signal') | |
| 13 | 13 | ||
@@ -79,77 +79,66 @@ class StreamHandler extends AsyncResource { | |||
| 79 | 79 | } | |
| 80 | 80 | ||
| 81 | 81 | onHeaders (statusCode, rawHeaders, resume, statusMessage) { | |
| 82 | - const { factory, opaque, context, callback } = this | ||
| 82 | + const { factory, opaque, context, callback, responseHeaders } = this | ||
| 83 | + | ||
| 84 | + const headers = responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 83 | 85 | ||
| 84 | 86 | if (statusCode < 200) { | |
| 85 | 87 | if (this.onInfo) { | |
| 86 | - const headers = this.responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 87 | 88 | this.onInfo({ statusCode, headers }) | |
| 88 | 89 | } | |
| 89 | 90 | return | |
| 90 | 91 | } | |
| 91 | 92 | ||
| 92 | 93 | this.factory = null | |
| 93 | - const headers = this.responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 94 | - const res = this.runInAsyncScope(factory, null, { | ||
| 95 | - statusCode, | ||
| 96 | - headers, | ||
| 97 | - opaque, | ||
| 98 | - context | ||
| 99 | - }) | ||
| 100 | 94 | ||
| 101 | - if (this.throwOnError && statusCode >= 400) { | ||
| 102 | - const headers = this.responseHeaders === 'raw' ? util.parseRawHeaders(rawHeaders) : util.parseHeaders(rawHeaders) | ||
| 103 | - const chunks = [] | ||
| 104 | - const pt = new PassThrough() | ||
| 105 | - pt | ||
| 106 | - .on('data', (chunk) => chunks.push(chunk)) | ||
| 107 | - .on('end', () => { | ||
| 108 | - const payload = Buffer.concat(chunks).toString('utf8') | ||
| 109 | - this.runInAsyncScope( | ||
| 110 | - callback, | ||
| 111 | - null, | ||
| 112 | - new ResponseStatusCodeError( | ||
| 113 | - `Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, | ||
| 114 | - statusCode, | ||
| 115 | - headers, | ||
| 116 | - payload | ||
| 117 | - ) | ||
| 118 | - ) | ||
| 119 | - }) | ||
| 120 | - .on('error', (err) => { | ||
| 121 | - this.onError(err) | ||
| 122 | - }) | ||
| 123 | - this.res = pt | ||
| 124 | - return | ||
| 125 | - } | ||
| 95 | + let res | ||
| 126 | 96 | ||
| 127 | - if ( | ||
| 128 | - !res || | ||
| 129 | - typeof res.write !== 'function' || | ||
| 130 | - typeof res.end !== 'function' || | ||
| 131 | - typeof res.on !== 'function' | ||
| 132 | - ) { | ||
| 133 | - throw new InvalidReturnValueError('expected Writable') | ||
| 134 | - } | ||
| 97 | + if (this.throwOnError && statusCode >= 400) { | ||
| 98 | + const parsedHeaders = responseHeaders === 'raw' ? util.parseHeaders(rawHeaders) : headers | ||
| 99 | + const contentType = parsedHeaders['content-type'] | ||
| 100 | + res = new PassThrough() | ||
| 135 | 101 | ||
| 136 | - res.on('drain', resume) | ||
| 137 | - // TODO: Avoid finished. It registers an unnecessary amount of listeners. | ||
| 138 | - finished(res, { readable: false }, (err) => { | ||
| 139 | - const { callback, res, opaque, trailers, abort } = this | ||
| 102 | + this.callback = null | ||
| 103 | + this.runInAsyncScope(getResolveErrorBodyCallback, null, | ||
| 104 | + { callback, body: res, contentType, statusCode, statusMessage, headers } | ||
| 105 | + ) | ||
| 106 | + } else { | ||
| 107 | + res = this.runInAsyncScope(factory, null, { | ||
| 108 | + statusCode, | ||
| 109 | + headers, | ||
| 110 | + opaque, | ||
| 111 | + context | ||
| 112 | + }) | ||
| 140 | 113 | ||
| 141 | - this.res = null | ||
| 142 | - if (err || !res.readable) { | ||
| 143 | - util.destroy(res, err) | ||
| 114 | + if ( | ||
| 115 | + !res || | ||
| 116 | + typeof res.write !== 'function' || | ||
| 117 | + typeof res.end !== 'function' || | ||
| 118 | + typeof res.on !== 'function' | ||
| 119 | + ) { | ||
| 120 | + throw new InvalidReturnValueError('expected Writable') | ||
| 144 | 121 | } | |
| 145 | 122 | ||
| 146 | - this.callback = null | ||
| 147 | - this.runInAsyncScope(callback, null, err || null, { opaque, trailers }) | ||
| 123 | + // TODO: Avoid finished. It registers an unnecessary amount of listeners. | ||
| 124 | + finished(res, { readable: false }, (err) => { | ||
| 125 | + const { callback, res, opaque, trailers, abort } = this | ||
| 148 | 126 | ||
| 149 | - if (err) { | ||
| 150 | - abort() | ||
| 151 | - } | ||
| 152 | - }) | ||
| 127 | + this.res = null | ||
| 128 | + if (err || !res.readable) { | ||
| 129 | + util.destroy(res, err) | ||
| 130 | + } | ||
| 131 | + | ||
| 132 | + this.callback = null | ||
| 133 | + this.runInAsyncScope(callback, null, err || null, { opaque, trailers }) | ||
| 134 | + | ||
| 135 | + if (err) { | ||
| 136 | + abort() | ||
| 137 | + } | ||
| 138 | + }) | ||
| 139 | + } | ||
| 140 | + | ||
| 141 | + res.on('drain', resume) | ||
| 153 | 142 | ||
| 154 | 143 | this.res = res | |
| 155 | 144 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -17,11 +17,16 @@ const kAbort = Symbol('abort') | |||
| 17 | 17 | const kContentType = Symbol('kContentType') | |
| 18 | 18 | ||
| 19 | 19 | module.exports = class BodyReadable extends Readable { | |
| 20 | - constructor (resume, abort, contentType = '') { | ||
| 20 | + constructor ({ | ||
| 21 | + resume, | ||
| 22 | + abort, | ||
| 23 | + contentType = '', | ||
| 24 | + highWaterMark = 64 * 1024 // Same as nodejs fs streams. | ||
| 25 | + }) { | ||
| 21 | 26 | super({ | |
| 22 | 27 | autoDestroy: true, | |
| 23 | 28 | read: resume, | |
| 24 | - highWaterMark: 64 * 1024 // Same as nodejs fs streams. | ||
| 29 | + highWaterMark | ||
| 25 | 30 | }) | |
| 26 | 31 | ||
| 27 | 32 | this._readableState.dataEmitted = false | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -0,0 +1,46 @@ | |||
| 1 | + const assert = require('assert') | ||
| 2 | + const { | ||
| 3 | + ResponseStatusCodeError | ||
| 4 | + } = require('../core/errors') | ||
| 5 | + const { toUSVString } = require('../core/util') | ||
| 6 | + | ||
| 7 | + async function getResolveErrorBodyCallback ({ callback, body, contentType, statusCode, statusMessage, headers }) { | ||
| 8 | + assert(body) | ||
| 9 | + | ||
| 10 | + let chunks = [] | ||
| 11 | + let limit = 0 | ||
| 12 | + | ||
| 13 | + for await (const chunk of body) { | ||
| 14 | + chunks.push(chunk) | ||
| 15 | + limit += chunk.length | ||
| 16 | + if (limit > 128 * 1024) { | ||
| 17 | + chunks = null | ||
| 18 | + break | ||
| 19 | + } | ||
| 20 | + } | ||
| 21 | + | ||
| 22 | + if (statusCode === 204 || !contentType || !chunks) { | ||
| 23 | + process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers)) | ||
| 24 | + return | ||
| 25 | + } | ||
| 26 | + | ||
| 27 | + try { | ||
| 28 | + if (contentType.startsWith('application/json')) { | ||
| 29 | + const payload = JSON.parse(toUSVString(Buffer.concat(chunks))) | ||
| 30 | + process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers, payload)) | ||
| 31 | + return | ||
| 32 | + } | ||
| 33 | + | ||
| 34 | + if (contentType.startsWith('text/')) { | ||
| 35 | + const payload = toUSVString(Buffer.concat(chunks)) | ||
| 36 | + process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers, payload)) | ||
| 37 | + return | ||
| 38 | + } | ||
| 39 | + } catch (err) { | ||
| 40 | + // Process in a fallback if error | ||
| 41 | + } | ||
| 42 | + | ||
| 43 | + process.nextTick(callback, new ResponseStatusCodeError(`Response status code ${statusCode}${statusMessage ? `: ${statusMessage}` : ''}`, statusCode, headers)) | ||
| 44 | + } | ||
| 45 | + | ||
| 46 | + module.exports = { getResolveErrorBodyCallback } | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments