From 408ac84ba10415171b1bf723be9a991eaa566c17 Mon Sep 17 00:00:00 2001 From: Aya Morisawa Date: Sun, 6 Sep 2015 04:19:16 +0900 Subject: [PATCH] Add kefir.d.ts --- kefir/kefir-tests.ts | 222 +++++++++++++++++++++++++++++++++++++++++++ kefir/kefir.d.ts | 176 ++++++++++++++++++++++++++++++++++ 2 files changed, 398 insertions(+) create mode 100644 kefir/kefir-tests.ts create mode 100644 kefir/kefir.d.ts diff --git a/kefir/kefir-tests.ts b/kefir/kefir-tests.ts new file mode 100644 index 000000000..c423a3dfe --- /dev/null +++ b/kefir/kefir-tests.ts @@ -0,0 +1,222 @@ +/// + +import * as Kefir from 'kefir'; +import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'kefir'; + +//Create a stream +{ + let stream01: Stream = Kefir.never(); + let stream02: Stream = Kefir.later(1000, 1); + let stream03: Stream = Kefir.interval(1000, 1); + let stream04: Stream = Kefir.sequentially(1000, [1, 2, 3]); + { + let start = +new Date(); + let stream05: Stream = Kefir.fromPoll(1000, () => +new Date() - start); + } + { + let start = +new Date(); + let stream06: Stream = Kefir.withInterval(1000, function(emitter) { + var time = +new Date() - start; + if (time < 4000) { + emitter.emit(time); + } else { + emitter.end(); + } + }); + } + let stream07: Stream = Kefir.fromCallback(callback => setTimeout(() => callback(1), 1000)); + let stream08: Stream = Kefir.fromNodeCallback(callback => setTimeout(() => callback(null, 1), 1000)); + let stream09: Stream = Kefir.fromEvents(document.body, 'click'); + let stream10: Stream = Kefir.stream(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 = Kefir.constant(1); + let property02: Property = Kefir.constantError(1); + let property03: Property = Kefir.fromPromise(new Promise(fulfill => fulfill(1))); +} + +// Convert observables +{ + let property: Property = Kefir.sequentially(100, [1, 2, 3]).toProperty(() => 0); + let stream: Stream = 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 = Kefir.sequentially(100, [1, 2, 3]).map(x => x + 1); + let observable02: Stream = Kefir.sequentially(100, [1, 2, 3]).filter(x => x > 1); + let observable03: Stream = Kefir.sequentially(100, [1, 2, 3]).take(2); + let observable04: Stream = Kefir.sequentially(100, [1, 2, 3]).takeWhile(x => x < 3); + let observable05: Stream = Kefir.sequentially(100, [1, 2, 3]).last(); + let observable06: Stream = Kefir.sequentially(100, [1, 2, 3]).skip(2); + let observable07: Stream = Kefir.sequentially(100, [1, 3, 2]).skipWhile(x => x < 3); + let observable08: Stream = Kefir.sequentially(100, [1, 2, 2, 3, 1]).skipDuplicates(); + let observable09: Stream = Kefir.sequentially(100, [1, 2, 2.1, 3, 1]).skipDuplicates((a, b) => Math.round(a) === Math.round(b)); + let observable10: Stream = Kefir.sequentially(100, [1, 2, 2, 3]).diff((prev, next) => next - prev, 0); + let observable11: Stream = Kefir.sequentially(100, [1, 2, 2, 3]).scan((prev, next) => next + prev, 0); + let observable12: Stream = Kefir.sequentially(100, [[1], [], [2,3]]).flatten(); + let observable13: Stream = Kefir.sequentially(100, [1, 2, 3, 4]).flatten(x => x % 2 === 0 ? [x * 10] : []); + let observable14: Stream = Kefir.sequentially(200, [1, 2, 3]).delay(100); + let observable15: Stream = Kefir.sequentially(750, [1, 2, 3, 4, 5, 6, 7, 8, 9, 0]).throttle(2500); + let observable16: Stream = Kefir.sequentially(100, [1, 2, 3, 0, 0, 0, 4, 5, 6]).filter(x => x > 0).debounce(250); + let observable17: Stream = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => { + return {convert: x < 0, error: x * 2}; + }); + let observable18: Stream = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors().errorsToValues((x: number) => { + return {convert: x >= 0, value: x * 2}; + }); + let observable19: Stream = Kefir.sequentially(100, [0, 1, 2, 3]).valuesToErrors().mapErrors((x: number) => x * 2); + let observable20: Stream = Kefir.sequentially(100, [0, 1, 2, 3]).valuesToErrors().filterErrors((x: number) => (x % 2) === 0); + let observable21: Stream = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => { + return {convert: x < 0, error: x}; + }).endOnError(); + let observable22: Stream = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => { + return {convert: x < 0, error: x}; + }).skipValues(); + let observable23: Stream = Kefir.sequentially(100, [0, -1, 2, -3]).valuesToErrors(x => { + return {convert: x < 0, error: x}; + }).skipErrors(); + let observable24: Stream = Kefir.sequentially(100, [1, 2, 3]).skipEnd(); + let ovservable25: Stream = Kefir.sequentially(100, [1, 2, 3]).beforeEnd(() => 0); + let observable26: Stream = Kefir.sequentially(100, [1, 2, 3, 4, 5]).slidingWindow(3, 2) + let observable27: Stream = Kefir.sequentially(100, [1, 2, 3, 4, 5]).bufferWhile(x => x !== 3); + { + var myTransducer: any; + let observable28: Stream = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6]).transduce(myTransducer); + } + let observable28: Stream = Kefir.sequentially(100, [0, 1, 2, 3]).withHandler((emitter: Emitter, event: Event) => { + 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 = Kefir.sequentially(100, [1, 3]); + let b: Stream = Kefir.sequentially(100, [2, 4]).delay(40); + let observable01: Observable = Kefir.combine([a, b], (a, b) => a + b); + } + { + let a: Stream = Kefir.sequentially(100, [1, 3]); + let b: Stream = Kefir.sequentially(100, [2, 4]).delay(40); + let c: Stream = Kefir.sequentially(60, [5, 6, 7]); + let observable02: Observable = Kefir.combine([a, b], [c], (a, b, c) => a + b + c); + } + { + let a: Stream = Kefir.sequentially(100, [0, 1, 2, 3]); + let b: Stream = Kefir.sequentially(160, [4, 5, 6]); + let c: Property = Kefir.sequentially(100, [8, 9]).delay(260).toProperty(() => 7); + let observable03: Observable = Kefir.zip([a, b, c]); + } + { + let a: Stream = Kefir.sequentially(100, [0, 1, 2]); + let b: Stream = Kefir.sequentially(100, [0, 1, 2]).delay(30); + let c: Stream = Kefir.sequentially(100, [0, 1, 2]).delay(60); + let abc: Observable = Kefir.merge([a, b, c]); + } + { + let a: Stream = Kefir.sequentially(100, [0, 1, 2]); + let b: Stream = Kefir.sequentially(100, [3, 4, 5]); + let abc: Observable = Kefir.concat([a, b]); + } + { + let a: Stream = Kefir.sequentially(100, [0, 1, 2]); + let b: Stream = Kefir.sequentially(100, [0, 1, 2]).delay(30); + let c: Observable = Kefir.sequentially(100, [0, 1, 2]).delay(60); + let pool: ObservablePool = Kefir.pool(); + pool.plug(a); + pool.plug(b); + pool.plug(c); + } + let observable04: Observable = Kefir.repeat(i => { + if (i < 3) { + return Kefir.sequentially(100, [i, i]); + } else { + return false; + } + }); + let observable05: Stream = Kefir.sequentially(100, [1, 2, 3]).flatMap(x => Kefir.interval(40, x).take(4)); + let observable06: Stream = Kefir.sequentially(100, [1, 2, 3]).flatMapLatest(x => Kefir.interval(40, x).take(4)); + let observable07: Stream = Kefir.sequentially(100, [1, 2, 3]).flatMapFirst(x => Kefir.interval(40, x).take(4)); + let observable08: Stream = Kefir.sequentially(100, [1, 2, 3]).flatMapConcat(x => Kefir.interval(40, x).take(4)); + let observable09: Stream = Kefir.sequentially(100, [1, 2, 3]).flatMapConcurLimit(x => Kefir.interval(40, x).take(6), 2); + let observable10: Stream = Kefir.sequentially(100, [1, 2]).valuesToErrors().flatMapErrors(x => Kefir.interval(40, x).take(2)); +} + +// Combine two observables +{ + { + let foo: Stream = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6, 7, 8]); + let bar: Property = Kefir.sequentially(200, [false, true, false]).delay(40).toProperty(() => true); + let observable01: Stream = foo.filterBy(bar); + } + { + let a: Property = Kefir.sequentially(200, [2, 3]).toProperty(() => 1); + let b: Stream = Kefir.interval(100, 0).delay(40).take(5); + let observable02: Property = a.sampledBy(b) + } + { + let foo: Stream = Kefir.sequentially(100, [1, 2, 3, 4]); + let bar: Stream = Kefir.later(250, 0); + let observable03: Stream = foo.skipUntilBy(bar); + } + { + let foo: Stream = Kefir.sequentially(100, [1, 2, 3, 4]); + let bar: Stream = Kefir.later(250, 0); + let observable04: Stream = foo.takeUntilBy(bar); + } + { + let foo: Stream = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6, 7, 8]).delay(40); + let bar: Stream = Kefir.sequentially(300, [1, 2]) + let observable05: Stream = foo.bufferBy(bar); + } + { + let foo: Stream = Kefir.sequentially(100, [1, 2, 3, 4, 5, 6, 7, 8]); + let bar: Stream = Kefir.sequentially(200, [false, true, false]).delay(40); + let observable06: Stream = foo.bufferWhileBy(bar); + } + { + let foo: Stream = Kefir.sequentially(100, [1, 2, 3]); + let bar: Stream = Kefir.sequentially(100, [1, 2, 3]).delay(40); + let observable07: Stream = foo.awaiting(bar); + } +} diff --git a/kefir/kefir.d.ts b/kefir/kefir.d.ts new file mode 100644 index 000000000..6a9544a7c --- /dev/null +++ b/kefir/kefir.d.ts @@ -0,0 +1,176 @@ +// Type definitions for Kefir 2.8.1 +// Project: http://rpominov.github.io/kefir/ +// Definitions by: Aya Morisawa +// Definitions: https://github.com/borisyankov/DefinitelyTyped + +/// +/// + +declare module "kefir" { + export interface Observable { + // 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) => void): void; + offAny(callback: (event: Event) => void): void; + log(name?: string): void; + offLog(name?: string): void; + toPromise(PromiseConstructor?: typeof Promise): Promise; + } + + export interface Stream extends Observable { + toProperty(getCurrent?: () => T): Property; + + // Modify an stream + map(fn: (value: T) => U): Stream; + filter(predicate?: (value: T) => boolean): Stream; + take(n: number): Stream; + takeWhile(predicate?: (value: T) => boolean): Stream; + last(): Stream; + skip(n: number): Stream; + skipWhile(predicate?: (value: T) => boolean): Stream; + skipDuplicates(comparator?: (a: T, b: T) => boolean): Stream; + diff(fn?: (prev: T, next: T) => T, seed?: T): Stream; + scan(fn: (prev: T, next: T) => T, seed?: T): Stream; + flatten(transformer?: (value: T) => U[]): Stream; + delay(wait: number): Stream; + throttle(wait: number, options?: {leading: boolean, trailing: boolean}): Stream; + debounce(wait: number, options?: {immediate: boolean}): Stream; + valuesToErrors(handler?: (value: T) => {convert: boolean, error: U}): Stream; + errorsToValues(handler?: (error: S) => {convert: boolean, value: U}): Stream; + mapErrors(fn: (error: S) => U): Stream; + filterErrors(predicate?: (error: S) => boolean): Stream; + endOnError(): Stream; + skipValues(): Stream; + skipErrors(): Stream; + skipEnd(): Stream; + beforeEnd(fn: () => U): Stream; + slidingWindow(max: number, mix?: number): Stream; + bufferWhile(predicate: (value: T) => boolean): Stream; + transduce(transducer: any): Stream; + withHandler(handler: (emitter: Emitter, event: Event) => void): Stream; + + // Combine streams + combine(otherObs: Stream, combinator?: (value: T, ...values: U[]) => W): Stream; + zip(otherObs: Stream, combinator?: (value: T, ...values: U[]) => W): Stream; + merge(otherObs: Stream): Stream; + concat(otherObs: Stream): Stream; + flatMap(transform: (value: T) => Stream): Stream; + flatMapLatest(fn: (value: T) => Stream): Stream; + flatMapFirst(fn: (value: T) => Stream): Stream; + flatMapConcat(fn: (value: T) => Stream): Stream; + flatMapConcurLimit(fn: (value: T) => Stream, limit: number): Stream; + flatMapErrors(transform: (error: S) => Stream): Stream; + + // Combine two streams + filterBy(otherObs: Observable): Stream; + sampledBy(otherObs: Observable, combinator?: (a: T, b: U) => W): Stream; + skipUntilBy(otherObs: Observable): Stream; + takeUntilBy(otherObs: Observable): Stream; + bufferBy(otherObs: Observable, options?: {flushOnEnd: boolean}): Stream; + bufferWhileBy(otherObs: Observable): Stream; + awaiting(otherObs: Observable): Stream; + } + + export interface Property extends Observable { + changes(): Stream; + + // Modify an property + map(fn: (value: T) => U): Property; + filter(predicate?: (value: T) => boolean): Property; + take(n: number): Property; + takeWhile(predicate?: (value: T) => boolean): Property; + last(): Property; + skip(n: number): Property; + skipWhile(predicate?: (value: T) => boolean): Property; + skipDuplicates(comparator?: (a: T, b: T) => boolean): Property; + diff(fn?: (prev: T, next: T) => T, seed?: T): Property; + scan(fn: (prev: T, next: T) => T, seed?: T): Property; + flatten(transformer?: (value: T) => U[]): Property; + delay(wait: number): Property; + throttle(wait: number, options?: {leading: boolean, trailing: boolean}): Property; + debounce(wait: number, options?: {immediate: boolean}): Property; + valuesToErrors(handler?: (value: T) => {convert: boolean, error: U}): Property; + errorsToValues(handler?: (error: S) => {convert: boolean, value: U}): Property; + mapErrors(fn: (error: S) => U): Property; + filterErrors(predicate?: (error: S) => boolean): Property; + endOnError(): Property; + skipValues(): Property; + skipErrors(): Property; + skipEnd(): Property; + beforeEnd(fn: () => U): Property; + slidingWindow(max: number, mix?: number): Property; + bufferWhile(predicate: (value: T) => boolean): Property; + transduce(transducer: any): Property; + withHandler(handler: (emitter: Emitter, event: Event) => void): Property; + + // Combine properties + combine(otherObs: Property, combinator?: (value: T, ...values: U[]) => W): Property; + zip(otherObs: Property, combinator?: (value: T, ...values: U[]) => W): Property; + merge(otherObs: Property): Property; + concat(otherObs: Property): Property; + flatMap(transform: (value: T) => Property): Property; + flatMapLatest(fn: (value: T) => Property): Property; + flatMapFirst(fn: (value: T) => Property): Property; + flatMapConcat(fn: (value: T) => Property): Property; + flatMapConcurLimit(fn: (value: T) => Property, limit: number): Property; + flatMapErrors(transform: (error: S) => Property): Property; + + // Combine two properties + filterBy(otherObs: Observable): Property; + sampledBy(otherObs: Observable, combinator?: (a: T, b: U) => W): Property; + skipUntilBy(otherObs: Observable): Property; + takeUntilBy(otherObs: Observable): Property; + bufferBy(otherObs: Observable, options?: {flushOnEnd: boolean}): Property; + bufferWhileBy(otherObs: Observable): Property; + awaiting(otherObs: Observable): Property; + } + + export interface ObservablePool extends Observable { + plug(obs: Observable): void; + unPlug(obs: Observable): void; + } + + export interface Event { + type: string; + value: T; + current: boolean; + } + + export interface Emitter { + 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; + export function later(wait: number, value: T): Stream; + export function interval(interval: number, value: T): Stream; + export function sequentially(interval: number, values: T[]): Stream; + export function fromPoll(interval: number, fn: () => T): Stream; + export function withInterval(interval: number, handler: (emitter: Emitter) => void): Stream; + export function fromCallback(fn: (callback: (value: T) => void) => void): Stream; + export function fromNodeCallback(fn: (callback: (error: S, result: T) => void) => void): Stream; + export function fromEvents(target: EventTarget | NodeJS.EventEmitter | { on: Function, off: Function }, eventName: string, transform?: (value: T) => S): Stream; + export function stream(subscribe: (emitter: Emitter) => Function | void): Stream; + + // Create a property + export function constant(value: T): Property; + export function constantError(error: T): Property; + export function fromPromise(promise: Promise): Property; + + // Combine observables + export function combine(obss: Observable[], passiveObss: Observable[], combinator?: (...values: T[]) => U): Observable; + export function combine(obss: Observable[], combinator?: (...values: T[]) => U): Observable; + export function zip(obss: Observable[], passiveObss?: Observable[], combinator?: (...values: T[]) => U): Observable; + export function merge(obss: Observable[]): Observable; + export function concat(obss: Observable[]): Observable; + export function pool(): ObservablePool; + export function repeat(generator: (i: number) => Observable | boolean): Observable; +}