const {timer} = rxjs;
const {buffer, exhaustMap, takeUntil} = rxjs.operators;
const {TestScheduler} = rxjs.testing;
const {expect} = chai;
const test = (testName, testFn) => {
try {
testFn();
console.log(`Test PASS "${testName}"`);
} catch (error) {
console.error(`Test FAIL "${testName}"`, error.message);
}
}
const createTestScheduler = () => new TestScheduler((actual, expected) => {
expect(actual).deep.equal(expected);
});
const createTestStream = (source$, blocker$) => {
return source$.pipe(
buffer(source$.pipe(
exhaustMap(() => timer(10).pipe(
takeUntil(blocker$)
))
))
);
}
const testStream = ({ marbles, values}) => {
const testScheduler = createTestScheduler();
testScheduler.run((helpers) => {
const { cold, hot, expectObservable } = helpers;
const source$ = hot(marbles.source);
const blocker$ = hot(marbles.blocker);
const result$ = createTestStream(source$, blocker$)
expectObservable(result$).toBe(marbles.result, values.result);
});
}
test('should buffer changes with 10ms delay', () => {
testStream({
marbles: {
source: ' ^-a-b 7ms ---c 9ms -----| ',
blocker: '^- 10ms --- 10ms -----| ',
result: ' -- 10ms i-- 10ms j----(k|)',
},
values: {
result: {
i: ['a', 'b'],
j: ['c'],
k: [],
},
}
});
});
test('should block buffer in progress and move values to next one', () => {
testStream({
marbles: {
source: ' ^-a-b 7ms ---c 9ms -----| ',
blocker: '^- 8ms e---- 10ms -----| ',
result: ' -- 10ms --- 10ms j----(k|)',
},
values: {
result: {
j: ['a', 'b', 'c'],
k: [],
},
}
});
});
<script src="https://cdnjs.cloudflare.com/ajax/libs/chai/4.1.2/chai.js"></script>
<script src="https://unpkg.com/rxjs@^7/dist/bundles/rxjs.umd.min.js"></script>