Merge pull request #5702 from AyaMorisawa/kefir

Add kefir.d.ts
This commit is contained in:
Horiuchi_H
2015-09-08 10:32:37 +09:00
2 changed files with 398 additions and 0 deletions
+222
View File
@@ -0,0 +1,222 @@
/// <reference path="./kefir.d.ts" />
import * as Kefir from 'kefir';
import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'kefir';
//Create a stream
{
let stream01: Stream<void, void> = Kefir.never();
let stream02: Stream<number, void> = Kefir.later(1000, 1);
let stream03: Stream<number, void> = Kefir.interval(1000, 1);
let stream04: Stream<number, void> = Kefir.sequentially(1000, [1, 2, 3]);
{
let start = +new Date();
let stream05: Stream<number, void> = Kefir.fromPoll(1000, () => +new Date() - start);
}
{
let start = +new Date();
let stream06: Stream<number, void> = Kefir.withInterval<number, void>(1000, function(emitter) {
var time = +new Date() - start;
if (time < 4000) {
emitter.emit(time);
} else {
emitter.end();
}
});
}
let stream07: Stream<number, void> = Kefir.fromCallback<number>(callback => setTimeout(() => callback(1), 1000));
let stream08: Stream<number, void> = Kefir.fromNodeCallback<number, void>(callback => setTimeout(() => callback(null, 1), 1000));
let stream09: Stream<MouseEvent, void> = Kefir.fromEvents<MouseEvent, void>(document.body, 'click');
let stream10: Stream<number, void> = Kefir.stream<number, void>(emitter => {
let count = 0;
emitter.emit(count);
let intervalId = setInterval(() => {
count++;
if (count < 4) {
emitter.emit(count);
} else {
emitter.end();
}
}, 1000);
return () => clearInterval(intervalId);
});
}
// Create a property
{
let property01: Property<number, void> = Kefir.constant(1);
let property02: Property<void, number> = Kefir.constantError(1);
let property03: Property<number, void> = Kefir.fromPromise<number, void>(new Promise<number>(fulfill => fulfill(1)));
}
// Convert observables
{
let property: Property<number, void> = Kefir.sequentially(100, [1, 2, 3]).toProperty(() => 0);
let stream: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).toProperty(() => 0).changes();
}
// Subscribe / add side effects
{
Kefir.sequentially(1000, [1, 2]).onValue(x => console.log('value:', x));
Kefir.sequentially(1000, [1, 2]).offValue(x => console.log('value:', x));
Kefir.sequentially(1000, [1, 2]).valuesToErrors().onValue(x => console.log('error:', x));
Kefir.sequentially(1000, [1, 2]).valuesToErrors().offValue(x => console.log('error:', x));
Kefir.sequentially(1000, [1, 2]).onEnd(() => console.log('stream ended'));
Kefir.sequentially(1000, [1, 2]).offEnd(() => console.log('stream ended'));
Kefir.sequentially(1000, [1, 2]).onAny(event => console.log('event:', event));
Kefir.sequentially(1000, [1, 2]).offAny(event => console.log('event:', event));
Kefir.sequentially(1000, [1, 2]).log('my stream');
Kefir.sequentially(1000, [1, 2]).offLog('my stream');
Kefir.sequentially(1000, [1, 2]).toPromise().then(x => console.log('fulfilled with:', x));
}
// Modify an observable
{
let observable01: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).map(x => x + 1);
let observable02: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).filter(x => x > 1);
let observable03: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).take(2);
let observable04: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).takeWhile(x => x < 3);
let observable05: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).last();
let observable06: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).skip(2);
let observable07: Stream<number, void> = Kefir.sequentially(100, [1, 3, 2]).skipWhile(x => x < 3);
let observable08: Stream<number, void> = Kefir.sequentially(100, [1, 2, 2, 3, 1]).skipDuplicates();
let observable09: Stream<number, void> = Kefir.sequentially(100, [1, 2, 2.1, 3, 1]).skipDuplicates((a, b) => Math.round(a) === Math.round(b));
let observable10: Stream<number, void> = Kefir.sequentially(100, [1, 2, 2, 3]).diff((prev, next) => next - prev, 0);
let observable11: Stream<number, void> = Kefir.sequentially(100, [1, 2, 2, 3]).scan((prev, next) => next + prev, 0);
let observable12: Stream<number, void> = Kefir.sequentially(100, [[1], [], [2,3]]).flatten<number>();
let observable13: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4]).flatten<number>(x => x % 2 === 0 ? [x * 10] : []);
let observable14: Stream<number, void> = Kefir.sequentially(200, [1, 2, 3]).delay(100);
let observable15: Stream<number, void> = Kefir.sequentially(750, [1, 2, 3, 4, 5, 6, 7, 8, 9, 0]).throttle(2500);
let observable16: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 0, 0, 0, 4, 5, 6]).filter(x => x > 0).debounce(250);
let observable17: Stream<void, number> = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors<number>(x => {
return {convert: x < 0, error: x * 2};
});
let observable18: Stream<number, void> = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors<number>().errorsToValues<number>((x: number) => {
return {convert: x >= 0, value: x * 2};
});
let observable19: Stream<void, number> = Kefir.sequentially(100, [0, 1, 2, 3]).valuesToErrors<number>().mapErrors((x: number) => x * 2);
let observable20: Stream<void, number> = Kefir.sequentially(100, [0, 1, 2, 3]).valuesToErrors<number>().filterErrors((x: number) => (x % 2) === 0);
let observable21: Stream<void, number> = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => {
return {convert: x < 0, error: x};
}).endOnError();
let observable22: Stream<void, number> = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => {
return {convert: x < 0, error: x};
}).skipValues();
let observable23: Stream<void, void> = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => {
return {convert: x < 0, error: x};
}).skipErrors();
let observable24: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).skipEnd();
let ovservable25: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).beforeEnd(() => 0);
let observable26: Stream<number[], void> = Kefir.sequentially(100, [1, 2, 3, 4, 5]).slidingWindow(3, 2)
let observable27: Stream<number[], void> = Kefir.sequentially(100, [1, 2, 3, 4, 5]).bufferWhile(x => x !== 3);
{
var myTransducer: any;
let observable28: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6]).transduce<number>(myTransducer);
}
let observable28: Stream<number | string, void> = Kefir.sequentially(100, [0, 1, 2, 3]).withHandler<number | string, void>((emitter: Emitter<string | number, void>, event: Event<number>) => {
if (event.type === 'end') {
emitter.emit('bye');
emitter.end();
}
if (event.type === 'value') {
for (var i = 0; i < event.value; i++) {
emitter.emit(event.value);
}
}
});
}
// Combine observables
{
{
let a: Stream<number, void> = Kefir.sequentially(100, [1, 3]);
let b: Stream<number, void> = Kefir.sequentially(100, [2, 4]).delay(40);
let observable01: Observable<number, void> = Kefir.combine<number, void, number>([a, b], (a, b) => a + b);
}
{
let a: Stream<number, void> = Kefir.sequentially(100, [1, 3]);
let b: Stream<number, void> = Kefir.sequentially(100, [2, 4]).delay(40);
let c: Stream<number, void> = Kefir.sequentially(60, [5, 6, 7]);
let observable02: Observable<number, void> = Kefir.combine<number, void, number>([a, b], [c], (a, b, c) => a + b + c);
}
{
let a: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2, 3]);
let b: Stream<number, void> = Kefir.sequentially(160, [4, 5, 6]);
let c: Property<number, void> = Kefir.sequentially(100, [8, 9]).delay(260).toProperty(() => 7);
let observable03: Observable<number, void> = Kefir.zip<number, void, number>([a, b, c]);
}
{
let a: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2]);
let b: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2]).delay(30);
let c: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2]).delay(60);
let abc: Observable<number, void> = Kefir.merge<number, void>([a, b, c]);
}
{
let a: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2]);
let b: Stream<number, void> = Kefir.sequentially(100, [3, 4, 5]);
let abc: Observable<number, void> = Kefir.concat<number, void>([a, b]);
}
{
let a: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2]);
let b: Stream<number, void> = Kefir.sequentially(100, [0, 1, 2]).delay(30);
let c: Observable<number, void> = Kefir.sequentially(100, [0, 1, 2]).delay(60);
let pool: ObservablePool<number, void> = Kefir.pool<number, void>();
pool.plug(a);
pool.plug(b);
pool.plug(c);
}
let observable04: Observable<number, void> = Kefir.repeat<number, void>(i => {
if (i < 3) {
return Kefir.sequentially(100, [i, i]);
} else {
return false;
}
});
let observable05: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).flatMap(x => Kefir.interval(40, x).take(4));
let observable06: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).flatMapLatest(x => Kefir.interval(40, x).take(4));
let observable07: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).flatMapFirst(x => Kefir.interval(40, x).take(4));
let observable08: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).flatMapConcat(x => Kefir.interval(40, x).take(4));
let observable09: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).flatMapConcurLimit(x => Kefir.interval(40, x).take(6), 2);
let observable10: Stream<number, void> = Kefir.sequentially(100, [1, 2]).valuesToErrors().flatMapErrors(x => Kefir.interval(40, x).take(2));
}
// Combine two observables
{
{
let foo: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6, 7, 8]);
let bar: Property<boolean, void> = Kefir.sequentially(200, [false, true, false]).delay(40).toProperty(() => true);
let observable01: Stream<number, void> = foo.filterBy<void>(bar);
}
{
let a: Property<number, void> = Kefir.sequentially(200, [2, 3]).toProperty(() => 1);
let b: Stream<number, void> = Kefir.interval(100, 0).delay(40).take(5);
let observable02: Property<number, void> = a.sampledBy<number, void, number>(b)
}
{
let foo: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4]);
let bar: Stream<number, void> = Kefir.later(250, 0);
let observable03: Stream<number, void> = foo.skipUntilBy<number, void>(bar);
}
{
let foo: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4]);
let bar: Stream<number, void> = Kefir.later(250, 0);
let observable04: Stream<number, void> = foo.takeUntilBy<number, void>(bar);
}
{
let foo: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6, 7, 8]).delay(40);
let bar: Stream<number, void> = Kefir.sequentially(300, [1, 2])
let observable05: Stream<number[], void> = foo.bufferBy<number, void>(bar);
}
{
let foo: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6, 7, 8]);
let bar: Stream<boolean, void> = Kefir.sequentially(200, [false, true, false]).delay(40);
let observable06: Stream<number[], void> = foo.bufferWhileBy<void>(bar);
}
{
let foo: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]);
let bar: Stream<number, void> = Kefir.sequentially(100, [1, 2, 3]).delay(40);
let observable07: Stream<boolean, void> = foo.awaiting<number, void>(bar);
}
}
+176
View File
@@ -0,0 +1,176 @@
// Type definitions for Kefir 2.8.1
// Project: http://rpominov.github.io/kefir/
// Definitions by: Aya Morisawa <https://github.com/AyaMorisawa>
// Definitions: https://github.com/borisyankov/DefinitelyTyped
/// <reference path="../node/node.d.ts" />
/// <reference path="../es6-promise/es6-promise.d.ts" />
declare module "kefir" {
export interface Observable<T, S> {
// Subscribe / add side effects
onValue(callback: (value: T) => void): void;
offValue(callback: (value: T) => void): void;
onError(callback: (error: S) => void): void;
offError(callback: (error: S) => void): void;
onEnd(callback: () => void): void;
offEnd(callback: () => void): void;
onAny(callback: (event: Event<T | S>) => void): void;
offAny(callback: (event: Event<T | S>) => void): void;
log(name?: string): void;
offLog(name?: string): void;
toPromise(PromiseConstructor?: typeof Promise): Promise<T>;
}
export interface Stream<T, S> extends Observable<T, S> {
toProperty(getCurrent?: () => T): Property<T, S>;
// Modify an stream
map<U>(fn: (value: T) => U): Stream<U, S>;
filter(predicate?: (value: T) => boolean): Stream<T, S>;
take(n: number): Stream<T, S>;
takeWhile(predicate?: (value: T) => boolean): Stream<T, S>;
last(): Stream<T, S>;
skip(n: number): Stream<T, S>;
skipWhile(predicate?: (value: T) => boolean): Stream<T, S>;
skipDuplicates(comparator?: (a: T, b: T) => boolean): Stream<T, S>;
diff(fn?: (prev: T, next: T) => T, seed?: T): Stream<T, S>;
scan(fn: (prev: T, next: T) => T, seed?: T): Stream<T, S>;
flatten<U>(transformer?: (value: T) => U[]): Stream<U, S>;
delay(wait: number): Stream<T, S>;
throttle(wait: number, options?: {leading: boolean, trailing: boolean}): Stream<T, S>;
debounce(wait: number, options?: {immediate: boolean}): Stream<T, S>;
valuesToErrors<U>(handler?: (value: T) => {convert: boolean, error: U}): Stream<void, S | U>;
errorsToValues<U>(handler?: (error: S) => {convert: boolean, value: U}): Stream<T | U, void>;
mapErrors<U>(fn: (error: S) => U): Stream<T, U>;
filterErrors(predicate?: (error: S) => boolean): Stream<T, S>;
endOnError(): Stream<T, S>;
skipValues(): Stream<void, S>;
skipErrors(): Stream<T, void>;
skipEnd(): Stream<T, S>;
beforeEnd<U>(fn: () => U): Stream<T | U, S>;
slidingWindow(max: number, mix?: number): Stream<T[], S>;
bufferWhile(predicate: (value: T) => boolean): Stream<T[], S>;
transduce<U>(transducer: any): Stream<U, S>;
withHandler<U, V>(handler: (emitter: Emitter<U, S>, event: Event<T | S>) => void): Stream<U, S>;
// Combine streams
combine<U, V, W>(otherObs: Stream<U, V>, combinator?: (value: T, ...values: U[]) => W): Stream<W, S | V>;
zip<U, V, W>(otherObs: Stream<U, V>, combinator?: (value: T, ...values: U[]) => W): Stream<W, S | V>;
merge<U, V>(otherObs: Stream<U, V>): Stream<T | U, S | V>;
concat<U, V>(otherObs: Stream<U, V>): Stream<T | U, S | V>;
flatMap<U, V>(transform: (value: T) => Stream<U, V>): Stream<U, V>;
flatMapLatest<U, V>(fn: (value: T) => Stream<U, V>): Stream<U, V>;
flatMapFirst<U, V>(fn: (value: T) => Stream<U, V>): Stream<U, V>;
flatMapConcat<U, V>(fn: (value: T) => Stream<U, V>): Stream<U, V>;
flatMapConcurLimit<U, V>(fn: (value: T) => Stream<U, V>, limit: number): Stream<U, V>;
flatMapErrors<U, V>(transform: (error: S) => Stream<U, V>): Stream<U, V>;
// Combine two streams
filterBy<U>(otherObs: Observable<boolean, U>): Stream<T, S>;
sampledBy<U, V, W>(otherObs: Observable<U, V>, combinator?: (a: T, b: U) => W): Stream<W, S>;
skipUntilBy<U, V>(otherObs: Observable<U, V>): Stream<U, V>;
takeUntilBy<U, V>(otherObs: Observable<U, V>): Stream<U, V>;
bufferBy<U, V>(otherObs: Observable<U, V>, options?: {flushOnEnd: boolean}): Stream<T[], S>;
bufferWhileBy<U>(otherObs: Observable<boolean, U>): Stream<T[], S>;
awaiting<U, V>(otherObs: Observable<U, V>): Stream<boolean, S>;
}
export interface Property<T, S> extends Observable<T, S> {
changes(): Stream<T, S>;
// Modify an property
map<U>(fn: (value: T) => U): Property<U, S>;
filter(predicate?: (value: T) => boolean): Property<T, S>;
take(n: number): Property<T, S>;
takeWhile(predicate?: (value: T) => boolean): Property<T, S>;
last(): Property<T, S>;
skip(n: number): Property<T, S>;
skipWhile(predicate?: (value: T) => boolean): Property<T, S>;
skipDuplicates(comparator?: (a: T, b: T) => boolean): Property<T, S>;
diff(fn?: (prev: T, next: T) => T, seed?: T): Property<T, S>;
scan(fn: (prev: T, next: T) => T, seed?: T): Property<T, S>;
flatten<U>(transformer?: (value: T) => U[]): Property<U, S>;
delay(wait: number): Property<T, S>;
throttle(wait: number, options?: {leading: boolean, trailing: boolean}): Property<T, S>;
debounce(wait: number, options?: {immediate: boolean}): Property<T, S>;
valuesToErrors<U>(handler?: (value: T) => {convert: boolean, error: U}): Property<void, S | U>;
errorsToValues<U>(handler?: (error: S) => {convert: boolean, value: U}): Property<T | U, void>;
mapErrors<U>(fn: (error: S) => U): Property<T, U>;
filterErrors(predicate?: (error: S) => boolean): Property<T, S>;
endOnError(): Property<T, S>;
skipValues(): Property<void, S>;
skipErrors(): Property<T, void>;
skipEnd(): Property<T, S>;
beforeEnd<U>(fn: () => U): Property<T | U, S>;
slidingWindow(max: number, mix?: number): Property<T[], S>;
bufferWhile(predicate: (value: T) => boolean): Property<T[], S>;
transduce<U>(transducer: any): Property<U, S>;
withHandler<U, V>(handler: (emitter: Emitter<T, S>, event: Event<T | S>) => void): Property<U, S>;
// Combine properties
combine<U, V, W>(otherObs: Property<U, V>, combinator?: (value: T, ...values: U[]) => W): Property<W, S | V>;
zip<U, V, W>(otherObs: Property<U, V>, combinator?: (value: T, ...values: U[]) => W): Property<W, S | V>;
merge<U, V>(otherObs: Property<U, V>): Property<T | U, S | V>;
concat<U, V>(otherObs: Property<U, V>): Property<T | U, S | V>;
flatMap<U, V>(transform: (value: T) => Property<U, V>): Property<U, V>;
flatMapLatest<U, V>(fn: (value: T) => Property<U, V>): Property<U, V>;
flatMapFirst<U, V>(fn: (value: T) => Property<U, V>): Property<U, V>;
flatMapConcat<U, V>(fn: (value: T) => Property<U, V>): Property<U, V>;
flatMapConcurLimit<U, V>(fn: (value: T) => Property<U, V>, limit: number): Property<U, V>;
flatMapErrors<U, V>(transform: (error: S) => Property<U, V>): Property<U, V>;
// Combine two properties
filterBy<U>(otherObs: Observable<boolean, U>): Property<T, S>;
sampledBy<U, V, W>(otherObs: Observable<U, V>, combinator?: (a: T, b: U) => W): Property<W, S>;
skipUntilBy<U, V>(otherObs: Observable<U, V>): Property<U, V>;
takeUntilBy<U, V>(otherObs: Observable<U, V>): Property<U, V>;
bufferBy<U, V>(otherObs: Observable<U, V>, options?: {flushOnEnd: boolean}): Property<T[], S>;
bufferWhileBy<U>(otherObs: Observable<boolean, U>): Property<T[], S>;
awaiting<U, V>(otherObs: Observable<U, V>): Property<boolean, S>;
}
export interface ObservablePool<T, S> extends Observable<T, S> {
plug(obs: Observable<T, S>): void;
unPlug(obs: Observable<T, S>): void;
}
export interface Event<T> {
type: string;
value: T;
current: boolean;
}
export interface Emitter<T, S> {
emit(value: T): void;
error(error: S): void;
end(): void;
emitEvent(event: {type: string, value: T | S}): void;
}
// Create a stream
export function never(): Stream<void, void>;
export function later<T>(wait: number, value: T): Stream<T, void>;
export function interval<T>(interval: number, value: T): Stream<T, void>;
export function sequentially<T>(interval: number, values: T[]): Stream<T, void>;
export function fromPoll<T>(interval: number, fn: () => T): Stream<T, void>;
export function withInterval<T, S>(interval: number, handler: (emitter: Emitter<T, S>) => void): Stream<T, S>;
export function fromCallback<T>(fn: (callback: (value: T) => void) => void): Stream<T, void>;
export function fromNodeCallback<T, S>(fn: (callback: (error: S, result: T) => void) => void): Stream<T, S>;
export function fromEvents<T, S>(target: EventTarget | NodeJS.EventEmitter | { on: Function, off: Function }, eventName: string, transform?: (value: T) => S): Stream<T, S>;
export function stream<T, S>(subscribe: (emitter: Emitter<T, S>) => Function | void): Stream<T, S>;
// Create a property
export function constant<T>(value: T): Property<T, void>;
export function constantError<T>(error: T): Property<void, T>;
export function fromPromise<T, S>(promise: Promise<T>): Property<T, S>;
// Combine observables
export function combine<T, S, U>(obss: Observable<T, S>[], passiveObss: Observable<T, S>[], combinator?: (...values: T[]) => U): Observable<U, S>;
export function combine<T, S, U>(obss: Observable<T, S>[], combinator?: (...values: T[]) => U): Observable<U, S>;
export function zip<T, S, U>(obss: Observable<T, S>[], passiveObss?: Observable<T, S>[], combinator?: (...values: T[]) => U): Observable<U, S>;
export function merge<T, S>(obss: Observable<T, S>[]): Observable<T, S>;
export function concat<T, S>(obss: Observable<T, S>[]): Observable<T, S>;
export function pool<T, S>(): ObservablePool<T, S>;
export function repeat<T, S>(generator: (i: number) => Observable<T, S> | boolean): Observable<T, S>;
}