RxJS Playground
A blank canvas to play / experiment with RxJS: https://github.com/Reactive-Extensions/RxJS
HTML
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.0.1/Rx.min.js"></script>
JavaScript
Rx.Observable.prototype.combineAllCont = function(){
let lastValues = [];
return Rx.Observable.create(obs => {
let i = 0;
let subscription = this.subscribe(stream => {
const streamIndex = i;
subscription.add(stream.subscribe(res => {
lastValues[streamIndex] = res;
obs.next(lastValues);
}));
i++;
});
return subscription;
});
}
let searchBox = [
'',
'car',
'car ',
'car food',
'car food sun',
'cat food sun',
'cat sun'
]
function performSearch(word) {
return Rx.Observable.defer(() => {
console.log('sending search for ' + word);
return Rx.Observable.of({
word: word,
results: word.split('')
});
})
}
let searchBoxStream = Rx.Observable.interval(500)
.take(searchBox.length * 2)
.map((i) => searchBox[i % searchBox.length]);
let wordsStream = searchBoxStream
.map((str) => str.trim().split(' ').filter((w) => w.trim() != ''));
let wordSearchSubjects = [];
let wordSearchStreamSubject = new Rx.ReplaySubject(1);
wordsStream.subscribe((words) => {
const nWords = words.length;
const nSubjects = wordSearchSubjects.length;
// Update streams
for(i=0; i<nWords && i<nSubjects; i++) {
wordSearchSubjects[i].next(words[i]);
}
// Create streams
for(i=nSubjects; i<nWords; i++) {
const wordSearchSubject = new Rx.ReplaySubject(1);
wordSearchSubjects.push(wordSearchSubject);
wordSearchStreamSubject.next(
wordSearchSubject.asObservable()
.distinctUntilChanged()
.flatMap((w) => performSearch(w))
.concat(Rx.Observable.of(false)) // Ending signal
)
wordSearchSubjects[i].next(words[i]);
}
// Delete streams
for(i=nWords; i<nSubjects; i++) {
wordSearchSubjects[i].complete();
}
wordSearchSubjects.length = nWords;
});
let wordSearchStream = wordSearchStreamSubject
.combineAllCont()
.map((arr) => arr.filter((r) => r !== false));
resultingStream = wordSearchStream
.map((arr) => {
let ret = [];
...