import { Operator } from '../Operator'; import { Subscriber } from '../Subscriber'; import { Observable } from '../Observable'; import { OperatorFunction } from '../types'; import { SimpleOuterSubscriber, innerSubscribe, SimpleInnerSubscriber } from '../innerSubscribe'; /** * Buffers the source Observable values until `closingNotifier` emits. * * Collects values from the past as an array, and emits * that array only when another Observable emits. * * ![](buffer.png) * * Buffers the incoming Observable values until the given `closingNotifier` * Observable emits a value, at which point it emits the buffer on the output * Observable and starts a new buffer internally, awaiting the next time * `closingNotifier` emits. * * ## Example * * On every click, emit array of most recent interval events * * ```ts * import { fromEvent, interval } from 'rxjs'; * import { buffer } from 'rxjs/operators'; * * const clicks = fromEvent(document, 'click'); * const intervalEvents = interval(1000); * const buffered = intervalEvents.pipe(buffer(clicks)); * buffered.subscribe(x => console.log(x)); * ``` * * @see {@link bufferCount} * @see {@link bufferTime} * @see {@link bufferToggle} * @see {@link bufferWhen} * @see {@link window} * * @param {Observable} closingNotifier An Observable that signals the * buffer to be emitted on the output Observable. * @return {Observable} An Observable of buffers, which are arrays of * values. * @method buffer * @owner Observable */ export function buffer(closingNotifier: Observable): OperatorFunction { return function bufferOperatorFunction(source: Observable) { return source.lift(new BufferOperator(closingNotifier)); }; } class BufferOperator implements Operator { constructor(private closingNotifier: Observable) { } call(subscriber: Subscriber, source: any): any { return source.subscribe(new BufferSubscriber(subscriber, this.closingNotifier)); } } /** * We need this JSDoc comment for affecting ESDoc. * @ignore * @extends {Ignored} */ class BufferSubscriber extends SimpleOuterSubscriber { private buffer: T[] = []; constructor(destination: Subscriber, closingNotifier: Observable) { super(destination); this.add(innerSubscribe(closingNotifier, new SimpleInnerSubscriber(this))); } protected _next(value: T) { this.buffer.push(value); } notifyNext(): void { const buffer = this.buffer; this.buffer = []; this.destination.next!(buffer); } }