feat(view): changed event emitters to be observables
This commit is contained in:
@@ -37,27 +37,50 @@ class ObservableWrapper {
|
||||
return s.listen(onNext, onError: onError, onDone: onComplete, cancelOnError: true);
|
||||
}
|
||||
|
||||
static StreamController createController() {
|
||||
return new StreamController.broadcast();
|
||||
static void callNext(EventEmitter emitter, value) {
|
||||
emitter.add(value);
|
||||
}
|
||||
|
||||
static Stream createObservable(StreamController controller) {
|
||||
return controller.stream;
|
||||
static void callThrow(EventEmitter emitter, error) {
|
||||
emitter.addError(error);
|
||||
}
|
||||
|
||||
static void callNext(StreamController controller, value) {
|
||||
controller.add(value);
|
||||
}
|
||||
|
||||
static void callThrow(StreamController controller, error) {
|
||||
controller.addError(error);
|
||||
}
|
||||
|
||||
static void callReturn(StreamController controller) {
|
||||
controller.close();
|
||||
static void callReturn(EventEmitter emitter) {
|
||||
emitter.close();
|
||||
}
|
||||
}
|
||||
|
||||
class EventEmitter extends Stream {
|
||||
StreamController<String> _controller;
|
||||
|
||||
EventEmitter() {
|
||||
_controller = new StreamController.broadcast();
|
||||
}
|
||||
|
||||
StreamSubscription listen(void onData(String line), {
|
||||
void onError(Error error),
|
||||
void onDone(),
|
||||
bool cancelOnError }) {
|
||||
return _controller.stream.listen(onData,
|
||||
onError: onError,
|
||||
onDone: onDone,
|
||||
cancelOnError: cancelOnError);
|
||||
}
|
||||
|
||||
void add(value) {
|
||||
_controller.add(value);
|
||||
}
|
||||
|
||||
void addError(error) {
|
||||
_controller.addError(error);
|
||||
}
|
||||
|
||||
void close() {
|
||||
_controller.close();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
class _Completer {
|
||||
final Completer c;
|
||||
|
||||
|
||||
@@ -53,6 +53,28 @@ export class PromiseWrapper {
|
||||
}
|
||||
}
|
||||
|
||||
export class ObservableWrapper {
|
||||
static subscribe(emitter:EventEmitter, onNext, onThrow = null, onReturn = null) {
|
||||
return emitter.observer({next: onNext, throw: onThrow, return: onReturn});
|
||||
}
|
||||
|
||||
static callNext(emitter:EventEmitter, value:any) {
|
||||
emitter.next(value);
|
||||
}
|
||||
|
||||
static callThrow(emitter:EventEmitter, error:any) {
|
||||
emitter.throw(error);
|
||||
}
|
||||
|
||||
static callReturn(emitter:EventEmitter) {
|
||||
emitter.return();
|
||||
}
|
||||
}
|
||||
|
||||
//TODO: vsavkin change to interface
|
||||
export class Observable {
|
||||
observer(generator:Function){}
|
||||
}
|
||||
|
||||
/**
|
||||
* Use Rx.Observable but provides an adapter to make it work as specified here:
|
||||
@@ -60,39 +82,37 @@ export class PromiseWrapper {
|
||||
*
|
||||
* Once a reference implementation of the spec is available, switch to it.
|
||||
*/
|
||||
export var Observable = Rx.Observable;
|
||||
export var ObservableController = Rx.Subject;
|
||||
export class EventEmitter extends Observable {
|
||||
_subject:Rx.Subject;
|
||||
|
||||
export class ObservableWrapper {
|
||||
static createController():Rx.Subject {
|
||||
return new Rx.Subject();
|
||||
constructor() {
|
||||
super();
|
||||
this._subject = new Rx.Subject();
|
||||
}
|
||||
|
||||
static createObservable(subject:Rx.Subject):Observable {
|
||||
return subject;
|
||||
observer(generator) {
|
||||
// Rx.Scheduler.immediate and setTimeout is a workaround, so Rx works with zones.js.
|
||||
// Once https://github.com/angular/zone.js/issues/51 is fixed, the hack should be removed.
|
||||
return this._subject.observeOn(Rx.Scheduler.immediate).subscribe(
|
||||
(value) => {setTimeout(() => generator.next(value));},
|
||||
(error) => generator.throw ? generator.throw(error) : null,
|
||||
() => generator.return ? generator.return() : null
|
||||
);
|
||||
}
|
||||
|
||||
static subscribe(observable:Observable, generatorOrOnNext, onThrow = null, onReturn = null) {
|
||||
if (isPresent(generatorOrOnNext.next)) {
|
||||
return observable.observeOn(Rx.Scheduler.timeout).subscribe(
|
||||
(value) => generatorOrOnNext.next(value),
|
||||
(error) => generatorOrOnNext.throw(error),
|
||||
() => generatorOrOnNext.return()
|
||||
);
|
||||
} else {
|
||||
return observable.observeOn(Rx.Scheduler.timeout).subscribe(generatorOrOnNext, onThrow, onReturn);
|
||||
}
|
||||
toRx():Rx.Observable {
|
||||
return this._subject;
|
||||
}
|
||||
|
||||
static callNext(subject:Rx.Subject, value:any) {
|
||||
subject.onNext(value);
|
||||
next(value) {
|
||||
this._subject.onNext(value);
|
||||
}
|
||||
|
||||
static callThrow(subject:Rx.Subject, error:any) {
|
||||
subject.onError(error);
|
||||
throw(error) {
|
||||
this._subject.onError(error);
|
||||
}
|
||||
|
||||
static callReturn(subject:Rx.Subject) {
|
||||
subject.onCompleted();
|
||||
return(value) {
|
||||
this._subject.onCompleted();
|
||||
}
|
||||
}
|
||||
@@ -49,35 +49,50 @@ export class PromiseWrapper {
|
||||
}
|
||||
|
||||
|
||||
export class ObservableWrapper {
|
||||
static subscribe(emitter: EventEmitter, onNext, onThrow = null, onReturn = null) {
|
||||
return emitter.observer({next: onNext, throw: onThrow, return: onReturn});
|
||||
}
|
||||
|
||||
static callNext(emitter: EventEmitter, value: any) { emitter.next(value); }
|
||||
|
||||
static callThrow(emitter: EventEmitter, error: any) { emitter.throw(error); }
|
||||
|
||||
static callReturn(emitter: EventEmitter) { emitter.return (null); }
|
||||
}
|
||||
|
||||
// TODO: vsavkin change to interface
|
||||
export class Observable {
|
||||
observer(generator: any) {}
|
||||
}
|
||||
|
||||
/**
|
||||
* Use Rx.Observable but provides an adapter to make it work as specified here:
|
||||
* https://github.com/jhusain/observable-spec
|
||||
*
|
||||
* Once a reference implementation of the spec is available, switch to it.
|
||||
*/
|
||||
type Observable = Rx.Observable<any>;
|
||||
type ObservableController = Rx.Subject<any>;
|
||||
export class EventEmitter extends Observable {
|
||||
_subject: Rx.Subject<any>;
|
||||
|
||||
export class ObservableWrapper {
|
||||
static createController(): Rx.Subject<any> { return new Rx.Subject(); }
|
||||
|
||||
static createObservable<T>(subject: Rx.Subject<T>): Rx.Observable<T> { return subject; }
|
||||
|
||||
static subscribe(observable: Rx.Observable<any>, generatorOrOnNext, onThrow = null,
|
||||
onReturn = null) {
|
||||
if (isPresent(generatorOrOnNext.next)) {
|
||||
return observable.observeOn(Rx.Scheduler.timeout)
|
||||
.subscribe((value) => generatorOrOnNext.next(value),
|
||||
(error) => generatorOrOnNext.throw(error), () => generatorOrOnNext.return ());
|
||||
} else {
|
||||
return observable.observeOn(Rx.Scheduler.timeout)
|
||||
.subscribe(generatorOrOnNext, onThrow, onReturn);
|
||||
}
|
||||
constructor() {
|
||||
super();
|
||||
this._subject = new Rx.Subject<any>();
|
||||
}
|
||||
|
||||
static callNext(subject: Rx.Subject<any>, value: any) { subject.onNext(value); }
|
||||
observer(generator) {
|
||||
var immediateScheduler = (<any>Rx.Scheduler).immediate;
|
||||
return this._subject.observeOn(immediateScheduler)
|
||||
.subscribe((value) => { setTimeout(() => generator.next(value)); },
|
||||
(error) => generator.throw ? generator.throw(error) : null,
|
||||
() => generator.return ? generator.return () : null);
|
||||
}
|
||||
|
||||
static callThrow(subject: Rx.Subject<any>, error: any) { subject.onError(error); }
|
||||
toRx(): Rx.Observable<any> { return this._subject; }
|
||||
|
||||
static callReturn(subject: Rx.Subject<any>) { subject.onCompleted(); }
|
||||
}
|
||||
next(value) { this._subject.onNext(value); }
|
||||
|
||||
throw(error) { this._subject.onError(error); }
|
||||
|
||||
return (value) { this._subject.onCompleted(); }
|
||||
}
|
||||
Reference in New Issue
Block a user