source: frontend/node_modules/readable-stream/lib/internal/streams/async_iterator.js

Last change on this file was 9af201e, checked in by MBK <marija.karapandzova@…>, 12 days ago

Fix frontend appearance

  • Property mode set to 100644
File size: 6.3 KB
Line 
1'use strict';
2
3var _Object$setPrototypeO;
4function _defineProperty(obj, key, value) { key = _toPropertyKey(key); if (key in obj) { Object.defineProperty(obj, key, { value: value, enumerable: true, configurable: true, writable: true }); } else { obj[key] = value; } return obj; }
5function _toPropertyKey(arg) { var key = _toPrimitive(arg, "string"); return typeof key === "symbol" ? key : String(key); }
6function _toPrimitive(input, hint) { if (typeof input !== "object" || input === null) return input; var prim = input[Symbol.toPrimitive]; if (prim !== undefined) { var res = prim.call(input, hint || "default"); if (typeof res !== "object") return res; throw new TypeError("@@toPrimitive must return a primitive value."); } return (hint === "string" ? String : Number)(input); }
7var finished = require('./end-of-stream');
8var kLastResolve = Symbol('lastResolve');
9var kLastReject = Symbol('lastReject');
10var kError = Symbol('error');
11var kEnded = Symbol('ended');
12var kLastPromise = Symbol('lastPromise');
13var kHandlePromise = Symbol('handlePromise');
14var kStream = Symbol('stream');
15function createIterResult(value, done) {
16 return {
17 value: value,
18 done: done
19 };
20}
21function readAndResolve(iter) {
22 var resolve = iter[kLastResolve];
23 if (resolve !== null) {
24 var data = iter[kStream].read();
25 // we defer if data is null
26 // we can be expecting either 'end' or
27 // 'error'
28 if (data !== null) {
29 iter[kLastPromise] = null;
30 iter[kLastResolve] = null;
31 iter[kLastReject] = null;
32 resolve(createIterResult(data, false));
33 }
34 }
35}
36function onReadable(iter) {
37 // we wait for the next tick, because it might
38 // emit an error with process.nextTick
39 process.nextTick(readAndResolve, iter);
40}
41function wrapForNext(lastPromise, iter) {
42 return function (resolve, reject) {
43 lastPromise.then(function () {
44 if (iter[kEnded]) {
45 resolve(createIterResult(undefined, true));
46 return;
47 }
48 iter[kHandlePromise](resolve, reject);
49 }, reject);
50 };
51}
52var AsyncIteratorPrototype = Object.getPrototypeOf(function () {});
53var ReadableStreamAsyncIteratorPrototype = Object.setPrototypeOf((_Object$setPrototypeO = {
54 get stream() {
55 return this[kStream];
56 },
57 next: function next() {
58 var _this = this;
59 // if we have detected an error in the meanwhile
60 // reject straight away
61 var error = this[kError];
62 if (error !== null) {
63 return Promise.reject(error);
64 }
65 if (this[kEnded]) {
66 return Promise.resolve(createIterResult(undefined, true));
67 }
68 if (this[kStream].destroyed) {
69 // We need to defer via nextTick because if .destroy(err) is
70 // called, the error will be emitted via nextTick, and
71 // we cannot guarantee that there is no error lingering around
72 // waiting to be emitted.
73 return new Promise(function (resolve, reject) {
74 process.nextTick(function () {
75 if (_this[kError]) {
76 reject(_this[kError]);
77 } else {
78 resolve(createIterResult(undefined, true));
79 }
80 });
81 });
82 }
83
84 // if we have multiple next() calls
85 // we will wait for the previous Promise to finish
86 // this logic is optimized to support for await loops,
87 // where next() is only called once at a time
88 var lastPromise = this[kLastPromise];
89 var promise;
90 if (lastPromise) {
91 promise = new Promise(wrapForNext(lastPromise, this));
92 } else {
93 // fast path needed to support multiple this.push()
94 // without triggering the next() queue
95 var data = this[kStream].read();
96 if (data !== null) {
97 return Promise.resolve(createIterResult(data, false));
98 }
99 promise = new Promise(this[kHandlePromise]);
100 }
101 this[kLastPromise] = promise;
102 return promise;
103 }
104}, _defineProperty(_Object$setPrototypeO, Symbol.asyncIterator, function () {
105 return this;
106}), _defineProperty(_Object$setPrototypeO, "return", function _return() {
107 var _this2 = this;
108 // destroy(err, cb) is a private API
109 // we can guarantee we have that here, because we control the
110 // Readable class this is attached to
111 return new Promise(function (resolve, reject) {
112 _this2[kStream].destroy(null, function (err) {
113 if (err) {
114 reject(err);
115 return;
116 }
117 resolve(createIterResult(undefined, true));
118 });
119 });
120}), _Object$setPrototypeO), AsyncIteratorPrototype);
121var createReadableStreamAsyncIterator = function createReadableStreamAsyncIterator(stream) {
122 var _Object$create;
123 var iterator = Object.create(ReadableStreamAsyncIteratorPrototype, (_Object$create = {}, _defineProperty(_Object$create, kStream, {
124 value: stream,
125 writable: true
126 }), _defineProperty(_Object$create, kLastResolve, {
127 value: null,
128 writable: true
129 }), _defineProperty(_Object$create, kLastReject, {
130 value: null,
131 writable: true
132 }), _defineProperty(_Object$create, kError, {
133 value: null,
134 writable: true
135 }), _defineProperty(_Object$create, kEnded, {
136 value: stream._readableState.endEmitted,
137 writable: true
138 }), _defineProperty(_Object$create, kHandlePromise, {
139 value: function value(resolve, reject) {
140 var data = iterator[kStream].read();
141 if (data) {
142 iterator[kLastPromise] = null;
143 iterator[kLastResolve] = null;
144 iterator[kLastReject] = null;
145 resolve(createIterResult(data, false));
146 } else {
147 iterator[kLastResolve] = resolve;
148 iterator[kLastReject] = reject;
149 }
150 },
151 writable: true
152 }), _Object$create));
153 iterator[kLastPromise] = null;
154 finished(stream, function (err) {
155 if (err && err.code !== 'ERR_STREAM_PREMATURE_CLOSE') {
156 var reject = iterator[kLastReject];
157 // reject if we are waiting for data in the Promise
158 // returned by next() and store the error
159 if (reject !== null) {
160 iterator[kLastPromise] = null;
161 iterator[kLastResolve] = null;
162 iterator[kLastReject] = null;
163 reject(err);
164 }
165 iterator[kError] = err;
166 return;
167 }
168 var resolve = iterator[kLastResolve];
169 if (resolve !== null) {
170 iterator[kLastPromise] = null;
171 iterator[kLastResolve] = null;
172 iterator[kLastReject] = null;
173 resolve(createIterResult(undefined, true));
174 }
175 iterator[kEnded] = true;
176 });
177 stream.on('readable', onReadable.bind(null, iterator));
178 return iterator;
179};
180module.exports = createReadableStreamAsyncIterator;
Note: See TracBrowser for help on using the repository browser.