import {defer, EMPTY, type Observable, of, type OperatorFunction, switchMap, tap} from 'rxjs' export function bufferUntil( emitWhen: (currentBuffer: T[]) => boolean, ): OperatorFunction { return (source: Observable) => defer(() => { let buffer: T[] = [] // custom buffer return source.pipe( tap((v) => buffer.push(v)), // add values to buffer switchMap(() => (emitWhen(buffer) ? of(buffer) : EMPTY)), // emit the buffer when the condition is met tap(() => (buffer = [])), // clear the buffer ) }) }