diff --git a/rx.js/rx-tests.ts b/rx.js/rx-tests.ts
new file mode 100644
index 000000000..14388c49f
--- /dev/null
+++ b/rx.js/rx-tests.ts
@@ -0,0 +1,135 @@
+///
+
+// Disposable
+var d: Rx.IDisposable = new Rx.Disposable(() => { });
+d = Rx.Disposable.create(() => { });
+d = Rx.Disposable.empty;
+d.dispose();
+
+// CompositeDisposable
+var cd = new Rx.CompositeDisposable(d, d, d);
+d = cd;
+cd = new Rx.CompositeDisposable([d, d]);
+cd.add(d);
+cd.clear();
+var b: boolean = cd.contains(d);
+var da: Rx.IDisposable[] = cd.toArray();
+cd.remove(d);
+
+// SingleAssignmentDisposable
+var sad = new Rx.SingleAssignmentDisposable();
+d = sad;
+sad.setDisposable(d);
+d = sad.getDisposable();
+b = sad.isDisposed;
+
+// SerialDisposable
+var sd = new Rx.SerialDisposable();
+d = sd;
+sd.setDisposable(d);
+d = sd.getDisposable();
+b = sd.isDisposed;
+
+// RefCountDisposable
+var rcd = new Rx.RefCountDisposable(d);
+d = rcd;
+d = rcd.getDisposable();
+b = rcd.isDisposed;
+
+// IScheduler
+var s: Rx.IScheduler;
+var n: number = s.now();
+s = s.catch(ex => true);
+s = s.catchException(ex => true);
+d = s.schedule(() => { });
+d = s.scheduleWithState(1,
+ (sh, s) => sh.scheduleWithAbsoluteAndState(s + 1, 100,
+ (sh, s) => sh.scheduleWithRelativeAndState(s + 1, 200,
+ (sh, s) => sh.scheduleRecursiveWithState(s + 1,
+ (s, self) => self(s + 1)))));
+d = s.scheduleWithAbsolute(100, () => { });
+d = s.scheduleWithRelative(100, () => { });
+d = s.scheduleRecursive(self => self());
+d = s.scheduleRecursiveWithAbsolute(100, self => self(200));
+d = s.scheduleRecursiveWithRelative(100, self => self(200));
+d = s.schedulePeriodic(100, () => { });
+d = s.schedulePeriodicWithState('a', 100, s => s + 'b');
+
+// ICurrentThreadScheduler
+Rx.Scheduler.currentThread.scheduleRequired();
+Rx.Scheduler.currentThread.ensureTrampoline(() => Rx.Disposable.empty);
+
+// Observer
+var o: Rx.Observer = Rx.Observer.create();
+o = Rx.Observer.create(i => { });
+o = Rx.Observer.create(i => { }, err => { });
+o = Rx.Observer.create(i => { }, err => { }, () => { });
+o = Rx.Observer.fromNotifier(n => { });
+
+o.onNext(10);
+o.onError(new Error());
+o.onCompleted();
+
+// Observable static methods tests
+
+var ns: Rx.Observable = Rx.Observable.create(observer => { });
+var ss: Rx.Observable = Rx.Observable.create(observer => (() => { }));
+var bs: Rx.Observable = Rx.Observable.createWithDisposable(observer => Rx.Disposable.empty);
+
+ns = Rx.Observable.defer(() => ns);
+
+ss = Rx.Observable.empty();
+bs = Rx.Observable.empty(Rx.Scheduler.currentThread);
+
+ns = Rx.Observable.fromArray([0, 3, -7, 18]);
+ss = Rx.Observable.fromArray(['a', 'ab', 'abc'], Rx.Scheduler.timeout);
+
+ns = Rx.Observable.generate(0, i => i < 100, i => i + 1, i => i * i);
+ns = Rx.Observable.generate(0, i => i < 100, i => i + 1, i => i * i, Rx.Scheduler.timeout);
+
+bs = Rx.Observable.never();
+
+ns = Rx.Observable.range(0, 100);
+ns = Rx.Observable.range(0, 100, Rx.Scheduler.timeout);
+
+ns = Rx.Observable.repeat(0, 100);
+ns = Rx.Observable.repeat(0, 100, Rx.Scheduler.timeout);
+
+ss = Rx.Observable.return('a');
+ss = Rx.Observable.return('a', Rx.Scheduler.timeout);
+ss = Rx.Observable.returnValue('a');
+ss = Rx.Observable.returnValue('a', Rx.Scheduler.timeout);
+
+bs = Rx.Observable.throw(new Error("error"));
+bs = Rx.Observable.throw(new Error("error"), Rx.Scheduler.timeout);
+bs = Rx.Observable.throwException(new Error("error"));
+bs = Rx.Observable.throwException(new Error("error"), Rx.Scheduler.timeout);
+
+bs = Rx.Observable.using(() => d, d => Rx.Observable.return(true));
+
+ss = Rx.Observable.amb(ss, ss);
+//ss = Rx.Observable.amb([ss, ss]);
+
+ss = Rx.Observable.catch(ss, ss, ss);
+ss = Rx.Observable.catchException(ss, ss, ss);
+//ss = Rx.Observable.catch([ss, ss, ss]);
+//ss = Rx.Observable.catchException([ss, ss, ss]);
+
+ss = Rx.Observable.concat(ss, ss, ss);
+//ss = Rx.Observable.concat([ss, ss, ss]);
+
+ss = Rx.Observable.merge(ss, ss, ss);
+//ss = Rx.Observable.merge([ss, ss, ss]);
+ss = Rx.Observable.merge(s, ss, ss, ss);
+//ss = Rx.Observable.merge(s, [ss, ss, ss]);
+
+ss = Rx.Observable.onErrorResumeNext(ss, ss, ss);
+//ss = Rx.Observable.onErrorResumeNext([ss, ss, ss]);
+
+ns = Rx.Observable.zip(ss, bs, ns, (s, b, n) => s.length + (b ? 1 : 0) + n);
+
+
+// Observable instance methods
+
+var sss: Rx.Observable>;
+ss = sss.concatAll();
\ No newline at end of file
diff --git a/rx.js/rx.js.aggregates.d.ts b/rx.js/rx.js.aggregates.d.ts
index 93a304db4..1955bfe7f 100644
--- a/rx.js/rx.js.aggregates.d.ts
+++ b/rx.js/rx.js.aggregates.d.ts
@@ -10,44 +10,44 @@
// -> rx.aggregates.js
declare module Rx {
- interface IObservable {
- aggregate(accumulator: (acc: TAcc, value: T) => TAcc): IObservable;
- aggregate(seed: TAcc, accumulator: (acc: TAcc, value: T) => TAcc): IObservable;
+ export interface Observable {
+ aggregate(accumulator: (acc: TAcc, value: T) => TAcc): Observable;
+ aggregate(seed: TAcc, accumulator: (acc: TAcc, value: T) => TAcc): Observable;
- any(): IObservable;
- any(selector: (item: T) => boolean): IObservable;
+ any(): Observable;
+ any(selector: (item: T) => boolean): Observable;
- isEmpty(predicate?: (value: T) => boolean): IObservable;
- all(predicate?: (value: T) => boolean): IObservable;
- contains(value: T, comparer?: (value1: T, value2: T) => boolean): IObservable;
- count(predicate?: (item: T) => boolean): IObservable;
- sum(keySelector?: (item: T) => number): IObservable;
- minBy(keySelector: (item: T) => number, comparer?: (value1: T, value2: T) => number): IObservable;
- min(comparer?: (value1: T, value2: T) => number): IObservable;
- maxBy(keySelector: (item: T) => number, comparer?: (value1: T, value2: T) => number): IObservable;
- max(comparer?: (value1: T, value2: T) => number): IObservable;
- average(keySelector?: (item: T) => number): IObservable;
+ isEmpty(predicate?: (value: T) => boolean): Observable;
+ all(predicate?: (value: T) => boolean): Observable;
+ contains(value: T, comparer?: (value1: T, value2: T) => boolean): Observable;
+ count(predicate?: (item: T) => boolean): Observable;
+ sum(keySelector?: (item: T) => number): Observable;
+ minBy(keySelector: (item: T) => number, comparer?: (value1: T, value2: T) => number): Observable;
+ min(comparer?: (value1: T, value2: T) => number): Observable;
+ maxBy(keySelector: (item: T) => number, comparer?: (value1: T, value2: T) => number): Observable;
+ max(comparer?: (value1: T, value2: T) => number): Observable;
+ average(keySelector?: (item: T) => number): Observable;
- sequenceEqual(second: IObservable, comparer?: (value1: T, value2: T) => number): IObservable;
- elementAt(index: number): IObservable;
- elementAtOrDefault(index: number, defaultValue: T): IObservable;
+ sequenceEqual(second: Observable, comparer?: (value1: T, value2: T) => number): Observable;
+ elementAt(index: number): Observable;
+ elementAtOrDefault(index: number, defaultValue: T): Observable;
- single(): IObservable;
- single(predicate: (T) => boolean): IObservable;
- singleOrDefault(): IObservable;
- singleOrDefault(predicate: (T) => boolean): IObservable;
- singleOrDefault(predicate: (T) => boolean, defaultValue: T): IObservable;
+ single(): Observable;
+ single(predicate: (item: T) => boolean): Observable;
+ singleOrDefault(): Observable;
+ singleOrDefault(predicate: (item: T) => boolean): Observable;
+ singleOrDefault(predicate: (item: T) => boolean, defaultValue: T): Observable;
- first(): IObservable;
- first(predicate: (T) => boolean): IObservable;
- firstOrDefault(): IObservable;
- firstOrDefault(predicate: (T) => boolean): IObservable;
- firstOrDefault(predicate: (T) => boolean, defaultValue: T): IObservable;
+ first(): Observable;
+ first(predicate: (item: T) => boolean): Observable;
+ firstOrDefault(): Observable;
+ firstOrDefault(predicate: (item: T) => boolean): Observable;
+ firstOrDefault(predicate: (item: T) => boolean, defaultValue: T): Observable;
- last(): IObservable;
- last(predicate: (T) => boolean): IObservable;
- lastOrDefault(): IObservable;
- lastOrDefault(predicate: (T) => boolean): IObservable;
- lastOrDefault(predicate: (T) => boolean, defaultValue: T): IObservable;
+ last(): Observable;
+ last(predicate: (item: T) => boolean): Observable;
+ lastOrDefault(): Observable;
+ lastOrDefault(predicate: (item: T) => boolean): Observable;
+ lastOrDefault(predicate: (item: T) => boolean, defaultValue: T): Observable;
}
}
\ No newline at end of file
diff --git a/rx.js/rx.js.async.d.ts b/rx.js/rx.js.async.d.ts
index 98fb930f7..fa336ac7e 100644
--- a/rx.js/rx.js.async.d.ts
+++ b/rx.js/rx.js.async.d.ts
@@ -11,12 +11,12 @@
declare module Rx {
interface Observable {
- start(func: () => T, scheduler?: IScheduler, context?: any): IObservable;
- toAsync(func: Function, scheduler?: IScheduler, context?: any): (...arguments: any[]) => IObservable;
- fromCallback(func: (...arguments: any[]) => void, scheduler?: IScheduler, context?: any, selector?: (...arguments: T[])=>T): () => IObservable;
- fromNodeCallback(func: (...arguments: any[]) => void, scheduler?: IScheduler, context?: any, selector?: (...arguments: any[])=>T): (...arguments: any[]) => IObservable;
- fromEvent(element: any, eventName: string, selector?: (...arguments: any[])=>T): IObservable;
- fromEventPattern(addHandler: (handler: any)=> void, removeHandler: (handler: any)=> void, selector?: (...arguments: any[])=>T): IObservable;
- fromPromise(promise: any): IObservable;
+ start(func: () => T, scheduler?: IScheduler, context?: any): Observable;
+ toAsync(func: Function, scheduler?: IScheduler, context?: any): (...arguments: any[]) => Observable;
+ fromCallback(func: (...arguments: any[]) => void, scheduler?: IScheduler, context?: any, selector?: (...arguments: T[])=>T): () => Observable;
+ fromNodeCallback(func: (...arguments: any[]) => void, scheduler?: IScheduler, context?: any, selector?: (...arguments: any[])=>T): (...arguments: any[]) => Observable;
+ fromEvent(element: any, eventName: string, selector?: (...arguments: any[])=>T): Observable;
+ fromEventPattern(addHandler: (handler: any)=> void, removeHandler: (handler: any)=> void, selector?: (...arguments: any[])=>T): Observable;
+ fromPromise(promise: any): Observable;
}
}
\ No newline at end of file
diff --git a/rx.js/rx.js.binding.d.ts b/rx.js/rx.js.binding.d.ts
index 97a62d0f8..0921f30c5 100644
--- a/rx.js/rx.js.binding.d.ts
+++ b/rx.js/rx.js.binding.d.ts
@@ -27,25 +27,25 @@ declare module Rx {
new (initialValue: T): BehaviorSubject;
}
- interface ConnectableObservable extends IObservable{
+ interface ConnectableObservable extends Observable{
connect(): IDisposable;
- refCount(): IObservable;
+ refCount(): Observable;
}
var ConnectableObservable: {
new (): ConnectableObservable;
}
- interface IObservable {
+ interface Observable {
publish(): ConnectableObservable;
- publish(selector: (item: T) => IObservable): ConnectableObservable;
+ publish(selector: (item: T) => Observable): ConnectableObservable;
publishLast(): ConnectableObservable;
- publishLast(selector: (item: T) => IObservable): ConnectableObservable;
+ publishLast(selector: (item: T) => Observable): ConnectableObservable;
publishValue(initialValue: T): ConnectableObservable;
publishValue(selector: (item: T) => TResult, initialValue: TResult): ConnectableObservable;
- replay(selector?: (source: IObservable) => ReplaySubject, bufferSize?: number, window?: number, scheduler?: IScheduler): ReplaySubject;
+ replay(selector?: (source: Observable) => ReplaySubject, bufferSize?: number, window?: number, scheduler?: IScheduler): ReplaySubject;
}
diff --git a/rx.js/rx.js.coincidence.d.ts b/rx.js/rx.js.coincidence.d.ts
index 138e99236..7aa0ae65a 100644
--- a/rx.js/rx.js.coincidence.d.ts
+++ b/rx.js/rx.js.coincidence.d.ts
@@ -11,27 +11,27 @@
declare module Rx {
- interface IObservable {
+ interface Observable {
join(
- right: IObservable,
- leftDurationSelector: (leftItem: T) => IObservable,
- rightDurationSelector: (rightItem: T2) => IObservable,
- resultSelector: (leftItem: T, rightItem: T2) => TResult): IObservable;
+ right: Observable,
+ leftDurationSelector: (leftItem: T) => Observable,
+ rightDurationSelector: (rightItem: T2) => Observable,
+ resultSelector: (leftItem: T, rightItem: T2) => TResult): Observable;
groupJoin(
- right: IObservable,
- leftDurationSelector: (leftItem: T) => IObservable,
- rightDurationSelector: (rightItem: T2) => IObservable,
- resultSelector: (leftItem: T, rightItem: IObservable) => TResult): IObservable;
+ right: Observable,
+ leftDurationSelector: (leftItem: T) => Observable,
+ rightDurationSelector: (rightItem: T2) => Observable,
+ resultSelector: (leftItem: T, rightItem: Observable) => TResult): Observable;
// lack of documentation to complete the followings...
- buffer(bufferOpenings: IObservable,
- bufferClosingSelector: (opening: TBufferOpening) => IObservable): IObservable;
+ buffer(bufferOpenings: Observable,
+ bufferClosingSelector: (opening: TBufferOpening) => Observable): Observable;
- window(bufferOpenings: IObservable,
- bufferClosingSelector: (opening: TBufferOpening) => IObservable): IObservable;
+ window(bufferOpenings: Observable,
+ bufferClosingSelector: (opening: TBufferOpening) => Observable): Observable;
}
diff --git a/rx.js/rx.js.d.ts b/rx.js/rx.js.d.ts
index 0b9c03de7..7107afb1e 100644
--- a/rx.js/rx.js.d.ts
+++ b/rx.js/rx.js.d.ts
@@ -8,7 +8,7 @@ declare module Rx {
export module Internals {
function inherits(child: Function, parent: Function): Function;
function addProperties(obj: Object, ...sourcces: Object[]): void;
- function addRef(xs: IObservable, r: { getDisposable(): IDisposable; }): IObservable;
+ function addRef(xs: Observable, r: { getDisposable(): IDisposable; }): Observable;
}
//Collections
@@ -34,7 +34,7 @@ declare module Rx {
remove(item: IIndexedItem): boolean;
}
- interface IDisposable {
+ export interface IDisposable {
dispose(): void;
}
@@ -59,9 +59,6 @@ declare module Rx {
static create(action: () => void): IDisposable;
static empty: IDisposable;
- isDisposed: boolean;
- action: () => void;
-
dispose(): void;
}
@@ -111,27 +108,27 @@ declare module Rx {
invokeCore(): IDisposable;
}
- interface IScheduler {
+ export interface IScheduler {
now(): number;
catch(handler: (exception: any) => boolean): IScheduler;
catchException(handler: (exception: any) => boolean): IScheduler;
schedule(action: () => void): IDisposable;
- scheduleWithState(state: any, action: (scheduler: IScheduler, state: any) => IDisposable): IDisposable;
+ scheduleWithState(state: TState, action: (scheduler: IScheduler, state: TState) => IDisposable): IDisposable;
scheduleWithAbsolute(dueTime: number, action: () => void): IDisposable;
- scheduleWithAbsoluteAndState(state: any, dueTime: number, action: (scheduler: IScheduler, state: any) =>IDisposable): IDisposable;
+ scheduleWithAbsoluteAndState(state: TState, dueTime: number, action: (scheduler: IScheduler, state: TState) =>IDisposable): IDisposable;
scheduleWithRelative(dueTime: number, action: () => void): IDisposable;
- scheduleWithRelativeAndState(state: any, dueTime: number, action: (scheduler: IScheduler, state: any) =>IDisposable): IDisposable;
+ scheduleWithRelativeAndState(state: TState, dueTime: number, action: (scheduler: IScheduler, state: TState) =>IDisposable): IDisposable;
scheduleRecursive(action: (action: () =>void ) =>void ): IDisposable;
- scheduleRecursiveWithState(state: any, action: (state: any, action: (state: any) =>void ) =>void ): IDisposable;
+ scheduleRecursiveWithState(state: TState, action: (state: TState, action: (state: TState) =>void ) =>void ): IDisposable;
scheduleRecursiveWithAbsolute(dueTime: number, action: (action: (dueTime: number) => void) => void): IDisposable;
- scheduleRecursiveWithAbsoluteAndState(state: any, dueTime: number, action: (state: any, action: (state: any, dueTime: number) => void) => void): IDisposable;
+ scheduleRecursiveWithAbsoluteAndState(state: TState, dueTime: number, action: (state: TState, action: (state: TState, dueTime: number) => void) => void): IDisposable;
scheduleRecursiveWithRelative(dueTime: number, action: (action: (dueTime: number) =>void ) =>void ): IDisposable;
- scheduleRecursiveWithRelativeAndState(state: any, dueTime: number, action: (state: any, action: (state: any, dueTime: number) =>void ) =>void ): IDisposable;
+ scheduleRecursiveWithRelativeAndState(state: TState, dueTime: number, action: (state: TState, action: (state: TState, dueTime: number) =>void ) =>void ): IDisposable;
schedulePeriodic(period: number, action: () => void): IDisposable;
- schedulePeriodicWithState(state: any, period: number, action: (state: any) => any): IDisposable;
+ schedulePeriodicWithState(state: TState, period: number, action: (state: TState) => TState): IDisposable;
}
export class Scheduler implements IScheduler {
@@ -152,21 +149,21 @@ declare module Rx {
catchException(handler: (exception: any) => boolean): IScheduler;
schedule(action: () => void): IDisposable;
- scheduleWithState(state: any, action: (scheduler: IScheduler, state: any) => IDisposable): IDisposable;
+ scheduleWithState(state: TState, action: (scheduler: IScheduler, state: TState) => IDisposable): IDisposable;
scheduleWithAbsolute(dueTime: number, action: () => void): IDisposable;
- scheduleWithAbsoluteAndState(state: any, dueTime: number, action: (scheduler: IScheduler, state: any) => IDisposable): IDisposable;
+ scheduleWithAbsoluteAndState(state: TState, dueTime: number, action: (scheduler: IScheduler, state: TState) => IDisposable): IDisposable;
scheduleWithRelative(dueTime: number, action: () => void): IDisposable;
- scheduleWithRelativeAndState(state: any, dueTime: number, action: (scheduler: IScheduler, state: any) => IDisposable): IDisposable;
+ scheduleWithRelativeAndState(state: TState, dueTime: number, action: (scheduler: IScheduler, state: TState) => IDisposable): IDisposable;
scheduleRecursive(action: (action: () => void) => void): IDisposable;
- scheduleRecursiveWithState(state: any, action: (state: any, action: (state: any) => void) => void): IDisposable;
+ scheduleRecursiveWithState(state: TState, action: (state: TState, action: (state: TState) => void) => void): IDisposable;
scheduleRecursiveWithAbsolute(dueTime: number, action: (action: (dueTime: number) => void) => void): IDisposable;
- scheduleRecursiveWithAbsoluteAndState(state: any, dueTime: number, action: (state: any, action: (state: any, dueTime: number) => void) => void): IDisposable;
+ scheduleRecursiveWithAbsoluteAndState(state: TState, dueTime: number, action: (state: TState, action: (state: TState, dueTime: number) => void) => void): IDisposable;
scheduleRecursiveWithRelative(dueTime: number, action: (action: (dueTime: number) => void) => void): IDisposable;
- scheduleRecursiveWithRelativeAndState(state: any, dueTime: number, action: (state: any, action: (state: any, dueTime: number) => void) => void): IDisposable;
+ scheduleRecursiveWithRelativeAndState(state: TState, dueTime: number, action: (state: TState, action: (state: TState, dueTime: number) => void) => void): IDisposable;
schedulePeriodic(period: number, action: () => void): IDisposable;
- schedulePeriodicWithState(state: any, period: number, action: (state: any) => any): IDisposable;
+ schedulePeriodicWithState(state: TState, period: number, action: (state: TState) => TState): IDisposable;
}
// Current Thread IScheduler
@@ -176,267 +173,193 @@ declare module Rx {
}
// Notifications
- interface INotification {
- accept(observer: IObserver): void;
- accept(onNext: (value: T) =>void , onError?: (exception: any) =>void , onCompleted?: () =>void ): void;
- toObservable(scheduler?: IScheduler): IObservable;
+ export class Notification {
+ accept(observer: Observer): void;
+ accept(onNext: (value: T) => TResult, onError?: (exception: any) => TResult, onCompleted?: () => TResult): TResult;
+ toObservable(scheduler?: IScheduler): Observable;
hasValue: boolean;
- equals(other: INotification): boolean;
+ equals(other: Notification): boolean;
kind: string;
- value?: T;
- exception?: any;
- }
- export interface Notification {
- //abstract
- //function new (): INotification;
+ value: T;
+ exception: any;
- createOnNext(value: T): INotification;//ON
- createOnError(exception): INotification;//OE
- createOnCompleted(): INotification;//OC
- }
-
- var Notification: Notification;
-
- export module Internals {
- // Enumerator
- interface IEnumerator {
- moveNext(): boolean;
- getCurrent(): T;
- dispose(): void;
- }
- export interface Enumerator {
- (moveNext: () =>boolean, getCurrent: () => T, dispose: () =>void ): IEnumerator;
-
- create(moveNext: () =>boolean, getCurrent: () => T, dispose?: () =>void ): IEnumerator;
- }
-
- // Enumerable
- interface IEnumerable {
- getEnumerator(): IEnumerator;
-
- concat(): IObservable;
- catchException(): IObservable;
- }
- export interface Enumerable {
- (getEnumerator: () =>IEnumerator): IEnumerable;
-
- repeat(value: T, repeatCount?: number): IEnumerable;
- forEach(source: T[], selector?: (element: T, index: number) => T2): IEnumerable;
- forEach(source: { length: number; [index: number]: T; }, selector?: (element: T, index: number) => T2): IEnumerable;
- }
+ static createOnNext(value: T): Notification;
+ static createOnError(exception): Notification;
+ static createOnCompleted(): Notification;
}
// Observer
- interface IObserver {
+ export class Observer {
onNext(value: T): void;
onError(exception: any): void;
onCompleted(): void;
- toNotifier(): (notification: INotification) =>void;
- asObserver(): IObserver;
- checked(): ICheckedObserver