/** * @license Use of this source code is governed by an MIT-style license that * can be found in the LICENSE file at https://github.com/cartant/rxjs-etc */ import { ConnectableObservable, from, Observable, ObservableInput, OperatorFunction, Subscription, } from "rxjs"; import { publish, startWith, switchMap, withLatestFrom } from "rxjs/operators"; export function withLatestFromWhen( flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( o2: ObservableInput, flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( o2: ObservableInput, o3: ObservableInput, flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( o2: ObservableInput, o3: ObservableInput, o4: ObservableInput, flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( o2: ObservableInput, o3: ObservableInput, o4: ObservableInput, o5: ObservableInput, flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( o2: ObservableInput, o3: ObservableInput, o4: ObservableInput, o5: ObservableInput, o6: ObservableInput, flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( array: ObservableInput[], flushSelector: () => Observable ): OperatorFunction; export function withLatestFromWhen( ...observables: (ObservableInput | (() => Observable))[] ): OperatorFunction; export function withLatestFromWhen( ...args: (ObservableInput | (() => Observable))[] ): OperatorFunction { const flushSelector = args.pop() as () => Observable; const observables = args as ObservableInput[]; return (source) => new Observable((subscriber) => { const publishedSource = publish()(source) as ConnectableObservable; const publishedObservables = observables.map( (o) => from(o).pipe(publish()) as ConnectableObservable ); const subscription = new Subscription(); subscription.add( flushSelector() .pipe( startWith(undefined), switchMap(() => publishedSource.pipe(withLatestFrom(...publishedObservables)) ) ) .subscribe(subscriber) ); publishedObservables.forEach((p) => subscription.add(p.connect())); subscription.add(publishedSource.connect()); return subscription; }); }