(stream, iterator, observer)
| 556 | } |
| 557 | |
| 558 | const wrapAsyncIterator = (stream, iterator, observer) => { |
| 559 | if (!observer || !iterator || typeof iterator.next !== 'function') { |
| 560 | return iterator; |
| 561 | } |
| 562 | |
| 563 | return { |
| 564 | async next(...nextArgs) { |
| 565 | try { |
| 566 | const result = await iterator.next(...nextArgs); |
| 567 | if (result.done) { |
| 568 | observer.closeConnection(stream); |
| 569 | } else { |
| 570 | observer.processChunk(result.value, stream); |
| 571 | } |
| 572 | return result; |
| 573 | } catch (error) { |
| 574 | observer.reportError(error); |
| 575 | throw error; |
| 576 | } |
| 577 | }, |
| 578 | async return(...returnArgs) { |
| 579 | observer.closeConnection(stream); |
| 580 | if (typeof iterator.return === 'function') { |
| 581 | return iterator.return(...returnArgs); |
| 582 | } |
| 583 | return { done: true }; |
| 584 | }, |
| 585 | async throw(...throwArgs) { |
| 586 | const error = throwArgs[0] || new Error('ReadableStream iterator error'); |
| 587 | observer.reportError(error); |
| 588 | if (typeof iterator.throw === 'function') { |
| 589 | return iterator.throw(...throwArgs); |
| 590 | } |
| 591 | throw error; |
| 592 | }, |
| 593 | [Symbol.asyncIterator]() { |
| 594 | return this; |
| 595 | } |
| 596 | }; |
| 597 | }; |
| 598 | |
| 599 | const createObservedIterator = function(originalIteratorFactory, iteratorArgs) { |
| 600 | const observer = streamObserverMap.get(this); |
no outgoing calls
no test coverage detected