Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Showing
2 changed files
with
114 additions
and
0 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,77 @@ | ||
/** | ||
* Created by tushar.mathur on 24/10/16. | ||
*/ | ||
|
||
import {IObservable} from '../types/core/IObservable' | ||
import {IObserver} from '../types/core/IObserver' | ||
import {IScheduler} from '../types/IScheduler' | ||
import {ISubscription} from '../types/core/ISubscription' | ||
import {LinkedList, LinkedListNode} from '../lib/LinkedList' | ||
|
||
export class MulticastSubscription<T> implements ISubscription { | ||
closed = false | ||
private node = this.sharedObserver.add(this.observer, this.scheduler) | ||
|
||
constructor (private observer: IObserver<T>, | ||
private scheduler: IScheduler, | ||
private sharedObserver: SharedObserver<T>) { | ||
} | ||
|
||
unsubscribe (): void { | ||
this.closed = true | ||
this.sharedObserver.remove(this.node) | ||
} | ||
} | ||
|
||
export class SharedObserver<T> implements IObserver<T> { | ||
private observers = new LinkedList<IObserver<T>>() | ||
private subscription: ISubscription | ||
|
||
constructor (private source: IObservable<T>) { | ||
} | ||
|
||
add (observer: IObserver<T>, scheduler: IScheduler) { | ||
const node = this.observers.add(observer) | ||
if (this.observers.length === 1) { | ||
this.subscription = this.source.subscribe(this, scheduler) | ||
} | ||
return node | ||
} | ||
|
||
remove (node: LinkedListNode<IObserver<T>>) { | ||
this.observers.remove(node) | ||
if (this.observers.length === 0) { | ||
this.subscription.unsubscribe() | ||
} | ||
} | ||
|
||
next (val: T): void { | ||
this.observers.forEach(ob => ob.value.next(val)) | ||
} | ||
|
||
error (err: Error): void { | ||
this.observers.forEach(ob => ob.value.error(err)) | ||
} | ||
|
||
complete (): void { | ||
this.observers.forEach(ob => ob.value.complete()) | ||
} | ||
|
||
} | ||
|
||
export class Multicast<T> implements IObservable<T> { | ||
constructor (private source: IObservable<T>) { | ||
} | ||
|
||
private sharedObserver = new SharedObserver(this.source) | ||
|
||
|
||
subscribe (observer: IObserver<T>, scheduler: IScheduler): ISubscription { | ||
|
||
return new MulticastSubscription(observer, scheduler, this.sharedObserver) | ||
} | ||
} | ||
|
||
export function multicast<T> (source: IObservable<T>) { | ||
return new Multicast(source) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,37 @@ | ||
/** | ||
* Created by tushar.mathur on 24/10/16. | ||
*/ | ||
|
||
import test from 'ava' | ||
import {TestScheduler} from '../src/testing/TestScheduler' | ||
import {ReactiveEvents} from '../src/testing/ReactiveEvents' | ||
import {map} from '../src/operators/Map' | ||
import {TestObserver} from '../src/testing/TestObserver' | ||
import {multicast} from '../src/operators/Multicast' | ||
|
||
test(t => { | ||
let i = 0 | ||
const results0: Array<number> = [] | ||
const results1: Array<number> = [] | ||
const sh = TestScheduler.of() | ||
const ob0 = new TestObserver(sh) | ||
const ob1 = new TestObserver(sh) | ||
const t$ = multicast(map((x: {(): number}) => x(), sh.Hot([ | ||
ReactiveEvents.next(10, () => ++i), | ||
ReactiveEvents.next(20, () => ++i), | ||
ReactiveEvents.next(30, () => ++i), | ||
ReactiveEvents.next(40, () => ++i), | ||
ReactiveEvents.next(50, () => ++i) | ||
]))) | ||
t$.subscribe(ob0, sh) | ||
t$.subscribe(ob1, sh) | ||
sh.advanceBy(50) | ||
t.deepEqual(ob0, ob1) | ||
t.deepEqual(ob0.results, [ | ||
ReactiveEvents.next(10, 1), | ||
ReactiveEvents.next(20, 2), | ||
ReactiveEvents.next(30, 3), | ||
ReactiveEvents.next(40, 4), | ||
ReactiveEvents.next(50, 5) | ||
]) | ||
}) |