import { checkObservable } from "../_check" import { NOOP_SUBSCRIBER, Source } from "../_core" import { In, Out } from "../_interfaces" import { makeObservable } from "../_obs" import { Transaction } from "../_tx" import { curry2 } from "../_util" import { dispatcherOf, Observable } from "../Observable" import { Pipe, PipeSubscriber } from "./_base" import { JoinOperator } from "./_join" import { SVSource } from "./sample" import { toEventStream } from "./toEventStream" export const takeUntil: CurriedTakeUntil = curry2(_takeUntil) export interface CurriedTakeUntil { (trigger: Observable, observable: In): Out< ObsType, ValueType > (trigger: Observable): ( observable: In, ) => Out } function _takeUntil(trigger: Observable, observable: Observable): Observable { checkObservable(trigger) return makeObservable( observable, new TakeUntil(dispatcherOf(toEventStream(trigger)), dispatcherOf(observable)), ) } class TakeUntil extends JoinOperator implements PipeSubscriber { protected source!: SVSource constructor(vSrc: Source, sSrc: Source) { super(new SVSource(vSrc, sSrc, NOOP_SUBSCRIBER)) this.source.vDest = new Pipe(this) } public next(tx: Transaction, val: T): void { this.forkNext(tx, val) } public pipedNext(sender: Pipe, tx: Transaction, v: any): void { this.source.disposeValue() this.sink.end(tx) } public pipedError(sender: Pipe, tx: Transaction, err: Error): void { // trigger errors are ignored } public pipedEnd(sender: Pipe, tx: Transaction): void { this.source.disposeValue() } }