(stream, type)
| 9796 | return util.isDisturbed(self) || isLocked(self); |
| 9797 | } |
| 9798 | function consume(stream, type) { |
| 9799 | return __async(this, null, function* () { |
| 9800 | assert(!stream[kConsume]); |
| 9801 | return new Promise((resolve2, reject) => { |
| 9802 | var _a; |
| 9803 | if (isUnusable(stream)) { |
| 9804 | const rState = stream._readableState; |
| 9805 | if (rState.destroyed && rState.closeEmitted === false) { |
| 9806 | stream.on("error", (err) => { |
| 9807 | reject(err); |
| 9808 | }).on("close", () => { |
| 9809 | reject(new TypeError("unusable")); |
| 9810 | }); |
| 9811 | } else { |
| 9812 | reject((_a = rState.errored) != null ? _a : new TypeError("unusable")); |
| 9813 | } |
| 9814 | } else { |
| 9815 | queueMicrotask(() => { |
| 9816 | stream[kConsume] = { |
| 9817 | type, |
| 9818 | stream, |
| 9819 | resolve: resolve2, |
| 9820 | reject, |
| 9821 | length: 0, |
| 9822 | body: [] |
| 9823 | }; |
| 9824 | stream.on("error", function(err) { |
| 9825 | consumeFinish(this[kConsume], err); |
| 9826 | }).on("close", function() { |
| 9827 | if (this[kConsume].body !== null) { |
| 9828 | consumeFinish(this[kConsume], new RequestAbortedError()); |
| 9829 | } |
| 9830 | }); |
| 9831 | consumeStart(stream[kConsume]); |
| 9832 | }); |
| 9833 | } |
| 9834 | }); |
| 9835 | }); |
| 9836 | } |
| 9837 | function consumeStart(consume2) { |
| 9838 | if (consume2.body === null) { |
| 9839 | return; |
no test coverage detected