Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
readable.js386 linesDownload Raw Back to api
1// Ported from https://github.com/nodejs/undici/pull/907
2
3'use strict'
4
5const assert = require('node:assert')
6const { Readable } = require('node:stream')
7const { RequestAbortedError, NotSupportedError, InvalidArgumentError, AbortError } = require('../core/errors')
8const util = require('../core/util')
9const { ReadableStreamFrom } = require('../core/util')
10
11const kConsume = Symbol('kConsume')
12const kReading = Symbol('kReading')
13const kBody = Symbol('kBody')
14const kAbort = Symbol('kAbort')
15const kContentType = Symbol('kContentType')
16const kContentLength = Symbol('kContentLength')
17
18const noop = () => {}
19
20class BodyReadable extends Readable {
21  constructor ({
22    resume,
23    abort,
24    contentType = '',
25    contentLength,
26    highWaterMark = 64 * 1024 // Same as nodejs fs streams.
27  }) {
28    super({
29      autoDestroy: true,
30      read: resume,
31      highWaterMark
32    })
33
34    this._readableState.dataEmitted = false
35
36    this[kAbort] = abort
37    this[kConsume] = null
38    this[kBody] = null
39    this[kContentType] = contentType
40    this[kContentLength] = contentLength
41
42    // Is stream being consumed through Readable API?
43    // This is an optimization so that we avoid checking
44    // for 'data' and 'readable' listeners in the hot path
45    // inside push().
46    this[kReading] = false
47  }
48
49  destroy (err) {
50    if (!err && !this._readableState.endEmitted) {
51      err = new RequestAbortedError()
52    }
53
54    if (err) {
55      this[kAbort]()
56    }
57
58    return super.destroy(err)
59  }
60
61  _destroy (err, callback) {
62    // Workaround for Node "bug". If the stream is destroyed in same
63    // tick as it is created, then a user who is waiting for a
64    // promise (i.e micro tick) for installing a 'error' listener will
65    // never get a chance and will always encounter an unhandled exception.
66    if (!this[kReading]) {
67      setImmediate(() => {
68        callback(err)
69      })
70    } else {
71      callback(err)
72    }
73  }
74
75  on (ev, ...args) {
76    if (ev === 'data' || ev === 'readable') {
77      this[kReading] = true
78    }
79    return super.on(ev, ...args)
80  }
81
82  addListener (ev, ...args) {
83    return this.on(ev, ...args)
84  }
85
86  off (ev, ...args) {
87    const ret = super.off(ev, ...args)
88    if (ev === 'data' || ev === 'readable') {
89      this[kReading] = (
90        this.listenerCount('data') > 0 ||
91        this.listenerCount('readable') > 0
92      )
93    }
94    return ret
95  }
96
97  removeListener (ev, ...args) {
98    return this.off(ev, ...args)
99  }
100
101  push (chunk) {
102    if (this[kConsume] && chunk !== null) {
103      consumePush(this[kConsume], chunk)
104      return this[kReading] ? super.push(chunk) : true
105    }
106    return super.push(chunk)
107  }
108
109  // https://fetch.spec.whatwg.org/#dom-body-text
110  async text () {
111    return consume(this, 'text')
112  }
113
114  // https://fetch.spec.whatwg.org/#dom-body-json
115  async json () {
116    return consume(this, 'json')
117  }
118
119  // https://fetch.spec.whatwg.org/#dom-body-blob
120  async blob () {
121    return consume(this, 'blob')
122  }
123
124  // https://fetch.spec.whatwg.org/#dom-body-bytes
125  async bytes () {
126    return consume(this, 'bytes')
127  }
128
129  // https://fetch.spec.whatwg.org/#dom-body-arraybuffer
130  async arrayBuffer () {
131    return consume(this, 'arrayBuffer')
132  }
133
134  // https://fetch.spec.whatwg.org/#dom-body-formdata
135  async formData () {
136    // TODO: Implement.
137    throw new NotSupportedError()
138  }
139
140  // https://fetch.spec.whatwg.org/#dom-body-bodyused
141  get bodyUsed () {
142    return util.isDisturbed(this)
143  }
144
145  // https://fetch.spec.whatwg.org/#dom-body-body
146  get body () {
147    if (!this[kBody]) {
148      this[kBody] = ReadableStreamFrom(this)
149      if (this[kConsume]) {
150        // TODO: Is this the best way to force a lock?
151        this[kBody].getReader() // Ensure stream is locked.
152        assert(this[kBody].locked)
153      }
154    }
155    return this[kBody]
156  }
157
158  async dump (opts) {
159    let limit = Number.isFinite(opts?.limit) ? opts.limit : 128 * 1024
160    const signal = opts?.signal
161
162    if (signal != null && (typeof signal !== 'object' || !('aborted' in signal))) {
163      throw new InvalidArgumentError('signal must be an AbortSignal')
164    }
165
166    signal?.throwIfAborted()
167
168    if (this._readableState.closeEmitted) {
169      return null
170    }
171
172    return await new Promise((resolve, reject) => {
173      if (this[kContentLength] > limit) {
174        this.destroy(new AbortError())
175      }
176
177      const onAbort = () => {
178        this.destroy(signal.reason ?? new AbortError())
179      }
180      signal?.addEventListener('abort', onAbort)
181
182      this
183        .on('close', function () {
184          signal?.removeEventListener('abort', onAbort)
185          if (signal?.aborted) {
186            reject(signal.reason ?? new AbortError())
187          } else {
188            resolve(null)
189          }
190        })
191        .on('error', noop)
192        .on('data', function (chunk) {
193          limit -= chunk.length
194          if (limit <= 0) {
195            this.destroy()
196          }
197        })
198        .resume()
199    })
200  }
201}
202
203// https://streams.spec.whatwg.org/#readablestream-locked
204function isLocked (self) {
205  // Consume is an implicit lock.
206  return (self[kBody] && self[kBody].locked === true) || self[kConsume]
207}
208
209// https://fetch.spec.whatwg.org/#body-unusable
210function isUnusable (self) {
211  return util.isDisturbed(self) || isLocked(self)
212}
213
214async function consume (stream, type) {
215  assert(!stream[kConsume])
216
217  return new Promise((resolve, reject) => {
218    if (isUnusable(stream)) {
219      const rState = stream._readableState
220      if (rState.destroyed && rState.closeEmitted === false) {
221        stream
222          .on('error', err => {
223            reject(err)
224          })
225          .on('close', () => {
226            reject(new TypeError('unusable'))
227          })
228      } else {
229        reject(rState.errored ?? new TypeError('unusable'))
230      }
231    } else {
232      queueMicrotask(() => {
233        stream[kConsume] = {
234          type,
235          stream,
236          resolve,
237          reject,
238          length: 0,
239          body: []
240        }
241
242        stream
243          .on('error', function (err) {
244            consumeFinish(this[kConsume], err)
245          })
246          .on('close', function () {
247            if (this[kConsume].body !== null) {
248              consumeFinish(this[kConsume], new RequestAbortedError())
249            }
250          })
251
252        consumeStart(stream[kConsume])
253      })
254    }
255  })
256}
257
258function consumeStart (consume) {
259  if (consume.body === null) {
260    return
261  }
262
263  const { _readableState: state } = consume.stream
264
265  if (state.bufferIndex) {
266    const start = state.bufferIndex
267    const end = state.buffer.length
268    for (let n = start; n < end; n++) {
269      consumePush(consume, state.buffer[n])
270    }
271  } else {
272    for (const chunk of state.buffer) {
273      consumePush(consume, chunk)
274    }
275  }
276
277  if (state.endEmitted) {
278    consumeEnd(this[kConsume])
279  } else {
280    consume.stream.on('end', function () {
281      consumeEnd(this[kConsume])
282    })
283  }
284
285  consume.stream.resume()
286
287  while (consume.stream.read() != null) {
288    // Loop
289  }
290}
291
292/**
293 * @param {Buffer[]} chunks
294 * @param {number} length
295 */
296function chunksDecode (chunks, length) {
297  if (chunks.length === 0 || length === 0) {
298    return ''
299  }
300  const buffer = chunks.length === 1 ? chunks[0] : Buffer.concat(chunks, length)
301  const bufferLength = buffer.length
302
303  // Skip BOM.
304  const start =
305    bufferLength > 2 &&
306    buffer[0] === 0xef &&
307    buffer[1] === 0xbb &&
308    buffer[2] === 0xbf
309      ? 3
310      : 0
311  return buffer.utf8Slice(start, bufferLength)
312}
313
314/**
315 * @param {Buffer[]} chunks
316 * @param {number} length
317 * @returns {Uint8Array}
318 */
319function chunksConcat (chunks, length) {
320  if (chunks.length === 0 || length === 0) {
321    return new Uint8Array(0)
322  }
323  if (chunks.length === 1) {
324    // fast-path
325    return new Uint8Array(chunks[0])
326  }
327  const buffer = new Uint8Array(Buffer.allocUnsafeSlow(length).buffer)
328
329  let offset = 0
330  for (let i = 0; i < chunks.length; ++i) {
331    const chunk = chunks[i]
332    buffer.set(chunk, offset)
333    offset += chunk.length
334  }
335
336  return buffer
337}
338
339function consumeEnd (consume) {
340  const { type, body, resolve, stream, length } = consume
341
342  try {
343    if (type === 'text') {
344      resolve(chunksDecode(body, length))
345    } else if (type === 'json') {
346      resolve(JSON.parse(chunksDecode(body, length)))
347    } else if (type === 'arrayBuffer') {
348      resolve(chunksConcat(body, length).buffer)
349    } else if (type === 'blob') {
350      resolve(new Blob(body, { type: stream[kContentType] }))
351    } else if (type === 'bytes') {
352      resolve(chunksConcat(body, length))
353    }
354
355    consumeFinish(consume)
356  } catch (err) {
357    stream.destroy(err)
358  }
359}
360
361function consumePush (consume, chunk) {
362  consume.length += chunk.length
363  consume.body.push(chunk)
364}
365
366function consumeFinish (consume, err) {
367  if (consume.body === null) {
368    return
369  }
370
371  if (err) {
372    consume.reject(err)
373  } else {
374    consume.resolve()
375  }
376
377  consume.type = null
378  consume.stream = null
379  consume.resolve = null
380  consume.reject = null
381  consume.length = 0
382  consume.body = null
383}
384
385module.exports = { Readable: BodyReadable, chunksDecode }
386 
codekingpro/portable-devtools · Team Ai