|
| 1 | +import Operator from '../Operator'; |
| 2 | +import Observer from '../Observer'; |
| 3 | +import Observable from '../Observable'; |
| 4 | +import Subscriber from '../Subscriber'; |
| 5 | + |
| 6 | +import {MergeSubscriber, MergeInnerSubscriber} from './merge'; |
| 7 | +import EmptyObservable from '../observables/EmptyObservable'; |
| 8 | +import ScalarObservable from '../observables/ScalarObservable'; |
| 9 | + |
| 10 | +import tryCatch from '../util/tryCatch'; |
| 11 | +import {errorObject} from '../util/errorObject'; |
| 12 | + |
| 13 | +export default function expand<T>(project: (x: T, ix: number) => Observable<any>): Observable<any> { |
| 14 | + return this.lift(new ExpandOperator(project)); |
| 15 | +} |
| 16 | + |
| 17 | +export class ExpandOperator<T, R> extends Operator<T, R> { |
| 18 | + |
| 19 | + project: (x: T, ix: number) => Observable<any>; |
| 20 | + |
| 21 | + constructor(project: (x: T, ix: number) => Observable<any>) { |
| 22 | + super(); |
| 23 | + this.project = project; |
| 24 | + } |
| 25 | + |
| 26 | + call(observer: Observer<R>): Observer<T> { |
| 27 | + return new ExpandSubscriber(observer, this.project); |
| 28 | + } |
| 29 | +} |
| 30 | + |
| 31 | +export class ExpandSubscriber<T, R> extends MergeSubscriber<T, R> { |
| 32 | + |
| 33 | + project: (x: T, ix: number) => Observable<any>; |
| 34 | + |
| 35 | + constructor(destination: Observer<R>, |
| 36 | + project: (x: T, ix: number) => Observable<any>) { |
| 37 | + super(destination, Number.POSITIVE_INFINITY); |
| 38 | + this.project = project; |
| 39 | + } |
| 40 | + |
| 41 | + _project(value, index) { |
| 42 | + const observable = tryCatch(this.project).call(this, value, index); |
| 43 | + if (observable === errorObject) { |
| 44 | + this.error(errorObject.e); |
| 45 | + return null; |
| 46 | + } |
| 47 | + return observable; |
| 48 | + } |
| 49 | + |
| 50 | + _subscribeInner(observable, value, index) { |
| 51 | + if(observable instanceof ScalarObservable) { |
| 52 | + this.destination.next((<ScalarObservable<T>> observable).value); |
| 53 | + this._innerComplete(); |
| 54 | + this._next((<ScalarObservable<T>> observable).value); |
| 55 | + } else if(observable instanceof EmptyObservable) { |
| 56 | + this._innerComplete(); |
| 57 | + } else { |
| 58 | + return observable.subscribe(new ExpandInnerSubscriber(this)); |
| 59 | + } |
| 60 | + } |
| 61 | +} |
| 62 | + |
| 63 | +export class ExpandInnerSubscriber<T, R> extends MergeInnerSubscriber<T, R> { |
| 64 | + constructor(parent: ExpandSubscriber<T, R>) { |
| 65 | + super(parent); |
| 66 | + } |
| 67 | + _next(value) { |
| 68 | + this.destination.next(value); |
| 69 | + this.parent.next(value); |
| 70 | + } |
| 71 | +} |
0 commit comments