From b1c7626cba4f53f87c5fd03086f170d8cfe5f4ad Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 09:51:47 +0400 Subject: [PATCH 01/11] Rx.Internals renamed to internals; Added Rx.config module. --- rx.js/rx.d.ts | 8 ++++++-- rx.js/rx.virtualtime.d.ts | 2 +- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/rx.js/rx.d.ts b/rx.js/rx.d.ts index a0f27b1bb..60c7c8634 100644 --- a/rx.js/rx.d.ts +++ b/rx.js/rx.d.ts @@ -1,11 +1,11 @@ -// Type definitions for RxJS v2.2.13 +// Type definitions for RxJS v2.2.15 // Project: http://rx.codeplex.com/ // Definitions by: gsino // Definitions by: Igor Oleinikov // Definitions: https://github.com/borisyankov/DefinitelyTyped declare module Rx { - export module Internals { + export module internals { function isEqual(left: any, right: any): boolean; function inherits(child: Function, parent: Function): Function; function addProperties(obj: Object, ...sourcces: Object[]): void; @@ -46,6 +46,10 @@ declare module Rx { } } + export module config { + + } + export interface IDisposable { dispose(): void; } diff --git a/rx.js/rx.virtualtime.d.ts b/rx.js/rx.virtualtime.d.ts index 376057964..afa398264 100644 --- a/rx.js/rx.virtualtime.d.ts +++ b/rx.js/rx.virtualtime.d.ts @@ -27,7 +27,7 @@ declare module Rx { /* protected abstract */ toDateTimeOffset(duetime: TAbsolute): number; /* protected abstract */ toRelative(duetime: number): TRelative; - /* protected */ getNext(): Internals.ScheduledItem; + /* protected */ getNext(): internals.ScheduledItem; } export class HistoricalScheduler extends VirtualTimeScheduler { From 3e29e44ddb330a33a5592780ef172ae746797a51 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 10:31:19 +0400 Subject: [PATCH 02/11] Promises was moved from rx.async.d.ts to rx.d.ts. --- rx.js/rx.async.d.ts | 25 ------------------------- rx.js/rx.d.ts | 40 +++++++++++++++++++++++++++++++++++++++- 2 files changed, 39 insertions(+), 26 deletions(-) diff --git a/rx.js/rx.async.d.ts b/rx.js/rx.async.d.ts index ef289aa0f..4ac758257 100644 --- a/rx.js/rx.async.d.ts +++ b/rx.js/rx.async.d.ts @@ -83,30 +83,5 @@ declare module Rx { fromEvent(element: NodeList, eventName: string, selector?: (arguments: any[]) => T): Observable; fromEvent(element: Node, eventName: string, selector?: (arguments: any[]) => T): Observable; fromEventPattern(addHandler: (handler: Function) => void, removeHandler: (handler: Function) => void, selector?: (arguments: any[])=>T): Observable; - - fromPromise(promise: IPromise): Observable; - fromPromise(promise: any): Observable; - } - - interface Observable { - /** - * Converts an existing observable sequence to an ES6 Compatible Promise - * @example - * var promise = Rx.Observable.return(42).toPromise(RSVP.Promise); - * @param The constructor of the promise - * @returns An ES6 compatible promise with the last value from the observable sequence. - */ - toPromise>(promiseCtor: { new (resolver: (resolvePromise: (value: T) => void, rejectPromise: (reason: any) => void) => void): TPromise; }): TPromise; - toPromise(promiseCtor: { new (resolver: (resolvePromise: (value: T) => void, rejectPromise: (reason: any) => void) => void): IPromise; }): IPromise; - } - - /** - * Promise A+ - */ - export interface IPromise { - then(onFulfilled: (value: T) => IPromise, onRejected: (reason: any) => IPromise): IPromise; - then(onFulfilled: (value: T) => IPromise, onRejected?: (reason: any) => R): IPromise; - then(onFulfilled: (value: T) => R, onRejected: (reason: any) => IPromise): IPromise; - then(onFulfilled?: (value: T) => R, onRejected?: (reason: any) => R): IPromise; } } diff --git a/rx.js/rx.d.ts b/rx.js/rx.d.ts index 60c7c8634..2b9a68dc9 100644 --- a/rx.js/rx.d.ts +++ b/rx.js/rx.d.ts @@ -47,7 +47,7 @@ declare module Rx { } export module config { - + export var Promise: { new (resolve: (value: T) => void, reject: (reason: any) => void): IPromise; }; } export interface IDisposable { @@ -189,6 +189,16 @@ declare module Rx { static createOnCompleted(): Notification; } + /** + * Promise A+ + */ + export interface IPromise { + then(onFulfilled: (value: T) => IPromise, onRejected: (reason: any) => IPromise): IPromise; + then(onFulfilled: (value: T) => IPromise, onRejected?: (reason: any) => R): IPromise; + then(onFulfilled: (value: T) => R, onRejected: (reason: any) => IPromise): IPromise; + then(onFulfilled?: (value: T) => R, onRejected?: (reason: any) => R): IPromise; + } + // Observer export class Observer { onNext(value: T): void; @@ -283,6 +293,27 @@ declare module Rx { takeWhile(predicate: (value: T, index: number, source: Observable) => boolean, thisArg?: any): Observable; where(predicate: (value: T, index: number, source: Observable) => boolean, thisArg?: any): Observable; filter(predicate: (value: T, index: number, source: Observable) => boolean, thisArg?: any): Observable; // alias for where + + /** + * Converts an existing observable sequence to an ES6 Compatible Promise + * @example + * var promise = Rx.Observable.return(42).toPromise(RSVP.Promise); + * @param promiseCtor The constructor of the promise. + * @returns An ES6 compatible promise with the last value from the observable sequence. + */ + toPromise>(promiseCtor: { new (resolver: (resolvePromise: (value: T) => void, rejectPromise: (reason: any) => void) => void): TPromise; }): TPromise; + /** + * Converts an existing observable sequence to an ES6 Compatible Promise + * @example + * var promise = Rx.Observable.return(42).toPromise(RSVP.Promise); + * + * // With config + * Rx.config.Promise = RSVP.Promise; + * var promise = Rx.Observable.return(42).toPromise(); + * @param [promiseCtor] The constructor of the promise. If not provided, it looks for it in Rx.config.Promise. + * @returns An ES6 compatible promise with the last value from the observable sequence. + */ + toPromise(promiseCtor?: { new (resolver: (resolvePromise: (value: T) => void, rejectPromise: (reason: any) => void) => void): IPromise; }): IPromise; } interface ObservableStatic { @@ -325,6 +356,13 @@ declare module Rx { zip(source1: Observable, source2: Observable, source3: Observable, source4: Observable, source5: Observable, resultSelector: (item1: T1, item2: T2, item3: T3, item4: T4, item5: T5) => TResult): Observable; zipArray(...sources: Observable[]): Observable; zipArray(sources: Observable[]): Observable; + + /** + * Converts a Promise to an Observable sequence + * @param promise An ES6 Compliant promise. + * @returns An Observable sequence which wraps the existing promise success and failure. + */ + fromPromise(promise: IPromise): Observable; } export var Observable: ObservableStatic; From 920c290ae36adcea961c6ceac1d93e24e0c80fba Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 10:46:14 +0400 Subject: [PATCH 03/11] Fixed definition for Rx.config.Promise; Added test for promises; --- rx.js/rx.async-tests.ts | 4 ++++ rx.js/rx.d.ts | 2 +- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/rx.js/rx.async-tests.ts b/rx.js/rx.async-tests.ts index b71f8bd77..9403c7517 100644 --- a/rx.js/rx.async-tests.ts +++ b/rx.js/rx.async-tests.ts @@ -66,8 +66,12 @@ module Rx.Tests.Async { new(resolver: (resolvePromise: (value: T)=> void, rejectPromise: (reason: any)=> void)=> void): Rx.IPromise; }; + Rx.config.Promise = promiseImpl; + var p: IPromise = obsNum.toPromise(promiseImpl); + p = obsNum.toPromise(); + p = p.then(x=> x); p = p.then(x=> p); p = p.then(undefined, reason=> 10); diff --git a/rx.js/rx.d.ts b/rx.js/rx.d.ts index 2b9a68dc9..d0dc1eda3 100644 --- a/rx.js/rx.d.ts +++ b/rx.js/rx.d.ts @@ -47,7 +47,7 @@ declare module Rx { } export module config { - export var Promise: { new (resolve: (value: T) => void, reject: (reason: any) => void): IPromise; }; + export var Promise: { new (resolver: (resolvePromise: (value: T) => void, rejectPromise: (reason: any) => void) => void): IPromise; }; } export interface IDisposable { From 35e263318a7cdf56b9e5f99568a4d3cee58539a2 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 10:59:20 +0400 Subject: [PATCH 04/11] Added Observable.fromGenerator method. --- rx.js/rx.d.ts | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/rx.js/rx.d.ts b/rx.js/rx.d.ts index d0dc1eda3..6d7cec2dd 100644 --- a/rx.js/rx.d.ts +++ b/rx.js/rx.d.ts @@ -363,6 +363,19 @@ declare module Rx { * @returns An Observable sequence which wraps the existing promise success and failure. */ fromPromise(promise: IPromise): Observable; + + /** + * Converts a generator function to an observable sequence, using an optional scheduler to enumerate the generator. + * BUG: it have been defined for Observable instance. + * + * @example + * var res = Rx.Observable.fromGenerator(function* () { yield 42; }); + * var res = Rx.Observable.fromArray(function* () { yield 42; }, Rx.Scheduler.timeout); + * @param genFn Generator function. + * @param [scheduler] Scheduler to run the enumeration of the input sequence on. + * @returns The observable sequence whose elements are pulled from the given generator sequence. + */ + fromGenerator(genFn: () => { next(): { done: boolean; value?: T; }; }, scheduler?: IScheduler): Observable; } export var Observable: ObservableStatic; From c3d5120566a8c2ecaac438f122fbae21fb3ca655 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 11:38:53 +0400 Subject: [PATCH 05/11] Added overloads of concat and merge taking promises instead of observables. --- rx.js/rx.d.ts | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/rx.js/rx.d.ts b/rx.js/rx.d.ts index 6d7cec2dd..1ab2d3574 100644 --- a/rx.js/rx.d.ts +++ b/rx.js/rx.d.ts @@ -233,11 +233,14 @@ declare module Rx { combineLatest(second: Observable, third: Observable, fourth: Observable, fifth: Observable, resultSelector: (v1: T, v2: T2, v3: T3, v4: T4, v5: T5) => TResult): Observable; combineLatest(souces: Observable[], resultSelector: (firstValue: T, ...otherValues: TOther[]) => TResult): Observable; concat(...sources: Observable[]): Observable; + concat(...sources: IPromise[]): Observable; concat(sources: Observable[]): Observable; + concat(sources: IPromise[]): Observable; concatAll(): T; concatObservable(): T; // alias for concatAll - merge(maxConcurrent: number): Observable; + merge(maxConcurrent: number): T; merge(other: Observable): Observable; + merge(other: IPromise): Observable; mergeAll(): T; mergeObservable(): T; // alias for mergeAll onErrorResumeNext(second: Observable): Observable; @@ -342,11 +345,17 @@ declare module Rx { catch(...sources: Observable[]): Observable; catchException(...sources: Observable[]): Observable; // alias for catch concat(...sources: Observable[]): Observable; + concat(...sources: IPromise[]): Observable; concat(sources: Observable[]): Observable; + concat(sources: IPromise[]): Observable; merge(...sources: Observable[]): Observable; + merge(...sources: IPromise[]): Observable; merge(sources: Observable[]): Observable; + merge(sources: IPromise[]): Observable; merge(scheduler: IScheduler, ...sources: Observable[]): Observable; + merge(scheduler: IScheduler, ...sources: IPromise[]): Observable; merge(scheduler: IScheduler, sources: Observable[]): Observable; + merge(scheduler: IScheduler, sources: IPromise[]): Observable; onErrorResumeNext(...sources: Observable[]): Observable; onErrorResumeNext(sources: Observable[]): Observable; zip(first: Observable, sources: Observable[], resultSelector: (item1: T1, right: Observable) => TResult): Observable; From 9664806862ed7943683fbe749713d49398e65bd3 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 11:46:07 +0400 Subject: [PATCH 06/11] Added overloads of `selectMany`/`flatMap` taking promises instead of observables. --- rx.js/rx.d.ts | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/rx.js/rx.d.ts b/rx.js/rx.d.ts index 1ab2d3574..640fe8ff7 100644 --- a/rx.js/rx.d.ts +++ b/rx.js/rx.d.ts @@ -245,7 +245,8 @@ declare module Rx { mergeObservable(): T; // alias for mergeAll onErrorResumeNext(second: Observable): Observable; skipUntil(other: Observable): Observable; - switchLatest(): T; + switch(): T; + switchLatest(): T; // alias for switch takeUntil(other: Observable): Observable; zip(second: Observable, resultSelector: (v1: T, v2: T2) => TResult): Observable; zip(second: Observable, third: Observable, resultSelector: (v1: T, v2: T2, v3: T3) => TResult): Observable; @@ -285,11 +286,17 @@ declare module Rx { select(selector: (value: T, index: number, source: Observable) => TResult, thisArg?: any): Observable; map(selector: (value: T, index: number, source: Observable) => TResult, thisArg?: any): Observable; // alias for select selectMany(selector: (value: T) => Observable, resultSelector: (item: T, other: TOther) => TResult): Observable; + selectMany(selector: (value: T) => IPromise, resultSelector: (item: T, other: TOther) => TResult): Observable; selectMany(selector: (value: T) => Observable): Observable; + selectMany(selector: (value: T) => IPromise): Observable; selectMany(other: Observable): Observable; + selectMany(other: IPromise): Observable; flatMap(selector: (value: T) => Observable, resultSelector: (item: T, other: TOther) => TResult): Observable; // alias for selectMany + flatMap(selector: (value: T) => IPromise, resultSelector: (item: T, other: TOther) => TResult): Observable; // alias for selectMany flatMap(selector: (value: T) => Observable): Observable; // alias for selectMany + flatMap(selector: (value: T) => IPromise): Observable; // alias for selectMany flatMap(other: Observable): Observable; // alias for selectMany + flatMap(other: IPromise): Observable; // alias for selectMany skip(count: number): Observable; skipWhile(predicate: (value: T, index: number, source: Observable) => boolean, thisArg?: any): Observable; take(count: number, scheduler?: IScheduler): Observable; From 5c938e56b08df58771158c4f7ff02f7466a626eb Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 11:51:52 +0400 Subject: [PATCH 07/11] Rx-Async: version bump to 2.2.15; added `ObservableStatic.startAsync`. --- rx.js/rx.async.d.ts | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/rx.js/rx.async.d.ts b/rx.js/rx.async.d.ts index 4ac758257..70143ff7e 100644 --- a/rx.js/rx.async.d.ts +++ b/rx.js/rx.async.d.ts @@ -1,4 +1,4 @@ -// Type definitions for RxJS-Async package 2.2.11 +// Type definitions for RxJS-Async package 2.2.15 // Project: http://rx.codeplex.com/ // Definitions by: zoetrope // Definitions by: Igor Oleinikov @@ -10,6 +10,13 @@ declare module Rx { interface ObservableStatic { start(func: () => T, scheduler?: IScheduler, context?: any): Observable; + /** + * Invokes the asynchronous function, surfacing the result through an observable sequence. + * @param functionAsync Asynchronous function which returns a Promise to run. + * @returns An observable sequence exposing the function's result value, or an exception. + */ + startAsync(functionAsync: () => IPromise): Observable; + toAsync(func: () => TResult, scheduler?: IScheduler, context?: any): () => Observable; toAsync(func: (arg1: T1) => TResult, scheduler?: IScheduler, context?: any): (arg1: T1) => Observable; toAsync(func: (arg1?: T1) => TResult, scheduler?: IScheduler, context?: any): (arg1?: T1) => Observable; From 6c4fcd603514d5664cffd0c32d1025231613b8e9 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 11:53:56 +0400 Subject: [PATCH 08/11] Rx-Async: added test for `startAsync`. --- rx.js/rx.async-tests.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/rx.js/rx.async-tests.ts b/rx.js/rx.async-tests.ts index 9403c7517..9c3bce516 100644 --- a/rx.js/rx.async-tests.ts +++ b/rx.js/rx.async-tests.ts @@ -81,4 +81,8 @@ module Rx.Tests.Async { ps = p.then(x=> ""); ps = p.then(x=> ps); } + + function startAsync() { + var o: Rx.Observable = Rx.Observable.startAsync(() => >null); + } } \ No newline at end of file From 7ae46cf1fd3408fe88737c26e23be8fff2b6caa0 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 12:01:02 +0400 Subject: [PATCH 09/11] Fixed headers of RxJS-Async and RxJS-Time. --- rx.js/rx.async.d.ts | 2 +- rx.js/rx.time.d.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/rx.js/rx.async.d.ts b/rx.js/rx.async.d.ts index 70143ff7e..19406e9ae 100644 --- a/rx.js/rx.async.d.ts +++ b/rx.js/rx.async.d.ts @@ -1,4 +1,4 @@ -// Type definitions for RxJS-Async package 2.2.15 +// Type definitions for RxJS-Async v2.2.15 // Project: http://rx.codeplex.com/ // Definitions by: zoetrope // Definitions by: Igor Oleinikov diff --git a/rx.js/rx.time.d.ts b/rx.js/rx.time.d.ts index 013665517..449cdd7f9 100644 --- a/rx.js/rx.time.d.ts +++ b/rx.js/rx.time.d.ts @@ -1,4 +1,4 @@ -// Type definitions for RxJS "Aggregates" +// Type definitions for RxJS-Time v2.2.15 // Project: http://rx.codeplex.com/ // Definitions by: Carl de Billy // Definitions by: Igor Oleinikov From 3a272478dacb005e79e6cb9ec4f1f417c999ba69 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 12:22:51 +0400 Subject: [PATCH 10/11] Added new `RxJS-BackPressure` module defintion. Fixed .gitignore to exclude rx.js from ignoring. --- .gitignore | 3 +++ rx.js/rx.backpressure.d.ts | 43 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 46 insertions(+) create mode 100644 rx.js/rx.backpressure.d.ts diff --git a/.gitignore b/.gitignore index 8fc575c21..b09429865 100644 --- a/.gitignore +++ b/.gitignore @@ -34,4 +34,7 @@ Properties *.iml *.js.map +#rx.js +!rx.js + node_modules diff --git a/rx.js/rx.backpressure.d.ts b/rx.js/rx.backpressure.d.ts new file mode 100644 index 000000000..d09aa91e7 --- /dev/null +++ b/rx.js/rx.backpressure.d.ts @@ -0,0 +1,43 @@ +// Type definitions for RxJS-BackPressure v2.2.15 +// Project: http://rx.codeplex.com/ +// Definitions by: Igor Oleinikov +// Definitions: https://github.com/borisyankov/DefinitelyTyped + +/// + +declare module Rx { + export interface Observable { + /** + * Pauses the underlying observable sequence based upon the observable sequence which yields true/false. + * @example + * var pauser = new Rx.Subject(); + * var source = Rx.Observable.interval(100).pausable(pauser); + * @param pauser The observable sequence used to pause the underlying sequence. + * @returns The observable sequence which is paused based upon the pauser. + */ + pausable(pauser: Observable): Observable; + + /** + * Pauses the underlying observable sequence based upon the observable sequence which yields true/false, + * and yields the values that were buffered while paused. + * @example + * var pauser = new Rx.Subject(); + * var source = Rx.Observable.interval(100).pausableBuffered(pauser); + * @param pauser The observable sequence used to pause the underlying sequence. + * @returns The observable sequence which is paused based upon the pauser. + */ + pausableBuffered(pauser: Observable): Observable; + + /** + * Attaches a controller to the observable sequence with the ability to queue. + * @example + * var source = Rx.Observable.interval(100).controlled(); + * source.request(3); // Reads 3 values + */ + controlled(enableQueue?: boolean): ControlledObservable; + } + + export interface ControlledObservable extends Observable { + request(numberOfItems?: number): IDisposable; + } +} From 064bcb14a7a279ba8d9976b2f0cf98b240810ed9 Mon Sep 17 00:00:00 2001 From: Igor Oleinikov Date: Tue, 18 Mar 2014 12:27:51 +0400 Subject: [PATCH 11/11] RxJS-BackPressure: added tests. --- rx.js/rx.backpressure-tests.ts | 22 ++++++++++++++++++++++ 1 file changed, 22 insertions(+) create mode 100644 rx.js/rx.backpressure-tests.ts diff --git a/rx.js/rx.backpressure-tests.ts b/rx.js/rx.backpressure-tests.ts new file mode 100644 index 000000000..f036a1dc3 --- /dev/null +++ b/rx.js/rx.backpressure-tests.ts @@ -0,0 +1,22 @@ +// Tests for RxJS-BackPressure TypeScript definitions +// Tests by Igor Oleinikov + +/// +/// + +function testPausable() { + var o: Rx.Observable; + + var pauser = new Rx.Subject(); + + var p = o.pausable(pauser); + p = o.pausableBuffered(pauser); +} + +function testControlled() { + var o: Rx.Observable; + var c = o.controlled(); + + var d: Rx.IDisposable = c.request(); + d = c.request(5); +}