From e6873eba9f250c2d2026d8fd05804e866ef8ebeb Mon Sep 17 00:00:00 2001 From: mAAdhaTTah Date: Sat, 5 Jan 2019 22:02:48 -0500 Subject: [PATCH] Improve kefir types Add proper type refinement to the `filter` method. Combine the operators on the base Observable class. Add staticLand exports. Add default export. --- types/kefir/index.d.ts | 335 +++++++++++++++++++------------------ types/kefir/kefir-tests.ts | 11 +- 2 files changed, 177 insertions(+), 169 deletions(-) diff --git a/types/kefir/index.d.ts b/types/kefir/index.d.ts index e892e1b291..05fd7a7aab 100644 --- a/types/kefir/index.d.ts +++ b/types/kefir/index.d.ts @@ -1,9 +1,9 @@ -// Type definitions for Kefir 3.7.3 +// Type definitions for Kefir 3.8.0 // Project: http://rpominov.github.io/kefir/ // Definitions by: Aya Morisawa // Piotr Hitori Bosak // Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped -// TypeScript Version: 2.4 +// TypeScript Version: 2.7 /// @@ -11,7 +11,7 @@ export type ValueOfAnObservable> = T['']; export interface Subscription { unsubscribe(): void; - closed: boolean; // Actually, `readonly` but it's avaiable in tsc starting with 2.0.0 + readonly closed: boolean; } export interface Observer { @@ -20,21 +20,24 @@ export interface Observer { end?: () => void; } -export interface Subscription { - unsubscribe(): void; - closed: boolean; // Actually, `readonly` but it's avaiable in tsc starting with 2.0.0 +interface ESObserver { + start?: Function, + next?: (value: T) => any, + error?: (error: S) => any, + complete?: () => any, } -export interface Observer { - value?: (value: T) => void; - error?: (error: S) => void; - end?: () => void; +interface ESObservable { + subscribe(callbacks: ESObserver): { unsubscribe(): void }; } -export interface Observable { +export class Observable { '': T; // TypeScript hack to enable value unwrapping for combine/flatMap - toProperty(getCurrent?: () => T): Property; + toProperty(): Property; + toProperty(getCurrent?: () => T2): Property; + changes(): Observable; + // Subscribe / add side effects onValue(callback: (value: T) => void): this; offValue(callback: (value: T) => void): this; @@ -42,8 +45,8 @@ export interface Observable { offError(callback: (error: S) => void): this; onEnd(callback: () => void): this; offEnd(callback: () => void): this; - onAny(callback: (event: Event) => void): this; - offAny(callback: (event: Event) => void): this; + onAny(callback: (event: Event) => void): this; + offAny(callback: (event: Event) => void): this; log(name?: string): this; spy(name?: string): this; offLog(name?: string): this; @@ -51,7 +54,7 @@ export interface Observable { flatten(transformer?: (value: T) => U[]): Stream; toPromise(): Promise; toPromise>(PromiseConstructor: () => W): W; - toESObservable(): any; + toESObservable(): ESObservable; // This method is designed to replace all other methods for subscribing observe(params: Observer): Subscription; observe( @@ -61,172 +64,174 @@ export interface Observable { ): Subscription; setName(source: Observable, selfName: string): this; setName(selfName: string): this; -} -export interface Stream extends Observable { + thru(cb: (obs: Observable) => Observable): Observable; // 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 | W, next: T) => W): Stream; - scan(fn: (prev: W, next: T) => W, seed: W): Stream; - delay(wait: number): Stream; - throttle(wait: number, options?: { leading?: boolean, trailing?: boolean }): Stream; - debounce(wait: number, options?: { immediate: boolean }): Stream; - valuesToErrors(): 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; - takeErrors(n: number): Stream; - ignoreValues(): Stream; - ignoreErrors(): Stream; - ignoreEnd(): Stream; - beforeEnd(fn: () => U): Stream; - slidingWindow(max: number, mix?: number): Stream; - bufferWhile(predicate: (value: T) => boolean): Stream; - bufferWithCount(count: number, options?: { flushOnEnd: boolean }): Stream; - bufferWithTimeOrCount(interval: number, count: number, options?: { flushOnEnd: boolean }): Stream; - transduce(transducer: any): Stream; - withHandler(handler: (emitter: Emitter, event: Event) => void): Stream; + map(fn: (value: T) => U): Observable; + filter(fn: (value: T) => value is U): Observable + filter(predicate?: (value: T) => boolean): Observable; + take(n: number): Observable; + takeWhile(predicate?: (value: T) => boolean): Observable; + last(): Observable; + skip(n: number): Observable; + skipWhile(predicate?: (value: T) => boolean): Observable; + skipDuplicates(comparator?: (a: T, b: T) => boolean): Observable; + diff(fn?: (prev: T, next: T) => T, seed?: T): Observable; + scan(fn: (prev: T | W, next: T) => W): Observable; + scan(fn: (prev: W, next: T) => W, seed: W): Observable; + delay(wait: number): Observable; + throttle(wait: number, options?: { leading?: boolean, trailing?: boolean }): Observable; + debounce(wait: number, options?: { immediate: boolean }): Observable; + valuesToErrors(): Observable; + valuesToErrors(handler: (value: T) => { convert: boolean, error: U }): Observable; + errorsToValues(handler?: (error: S) => { convert: boolean, value: U }): Observable; + mapErrors(fn: (error: S) => U): Observable; + filterErrors(predicate?: (error: S) => boolean): Observable; + endOnError(): Observable; + takeErrors(n: number): Observable; + ignoreValues(): Observable; + ignoreErrors(): Observable; + ignoreEnd(): Observable; + beforeEnd(fn: () => U): Observable; + slidingWindow(max: number, mix?: number): Observable; + bufferWhile(predicate: (value: T) => boolean): Observable; + bufferWithCount(count: number, options?: { flushOnEnd: boolean }): Observable; + bufferWithTimeOrCount(interval: number, count: number, options?: { flushOnEnd: boolean }): Observable; + transduce(transducer: any): Observable; + withHandler(handler: (emitter: Emitter, event: Event) => void): Observable; // 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; - flatMap>(): Stream, any>; - flatMapLatest(fn: (value: T) => Stream): Stream; - flatMapLatest>(): Stream, any>; - flatMapFirst(fn: (value: T) => Stream): Stream; - flatMapFirst>(): Stream, any>; - flatMapConcat(fn: (value: T) => Stream): Stream; - flatMapConcat>(): Stream, any>; - flatMapConcurLimit(fn: (value: T) => Stream, limit: number): Stream; - flatMapErrors(transform: (error: S) => Stream): Stream; + combine(otherObs: Observable, combinator?: (value: T, ...values: U[]) => W): Observable; + zip(otherObs: Observable, combinator?: (value: T, ...values: U[]) => W): Observable; + merge(otherObs: Observable): Observable; + concat(otherObs: Observable): Observable; + flatMap(transform: (value: T) => Observable): Observable; + flatMap>(): Observable, any>; + flatMapLatest(fn: (value: T) => Observable): Observable; + flatMapLatest>(): Observable, any>; + flatMapFirst(fn: (value: T) => Observable): Observable; + flatMapFirst>(): Observable, any>; + flatMapConcat(fn: (value: T) => Observable): Observable; + flatMapConcat>(): Observable, any>; + flatMapConcurLimit(fn: (value: T) => Observable, limit: number): Observable; + flatMapErrors(transform: (error: S) => Observable): Observable; // 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, options?: { flushOnEnd?: boolean, flushOnChange?: boolean }): Stream; - awaiting(otherObs: Observable): Stream; + filterBy(otherObs: Observable): Observable; + sampledBy(otherObs: Observable): Observable; + sampledBy(otherObs: Observable, combinator: (a: T, b: U) => W): Observable; + skipUntilBy(otherObs: Observable): Observable; + takeUntilBy(otherObs: Observable): Observable; + bufferBy(otherObs: Observable, options?: { flushOnEnd: boolean }): Observable; + bufferWhileBy(otherObs: Observable, options?: { flushOnEnd?: boolean, flushOnChange?: boolean }): Observable; + awaiting(otherObs: Observable): Observable; } -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; - 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; - takeErrors(n: number): Stream; - ignoreValues(): Property; - ignoreErrors(): Property; - ignoreEnd(): Property; - beforeEnd(fn: () => U): Property; - slidingWindow(max: number, mix?: number): Property; - bufferWhile(predicate: (value: T) => boolean): Property; - bufferWithCount(count: number, options?: { flushOnEnd: boolean }): Property; - bufferWithTimeOrCount(interval: number, count: number, options?: { flushOnEnd: 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; - flatMap>(): Property, any>; - flatMapLatest(fn: (value: T) => Property): Property; - flatMapLatest>(): Property, any>; - flatMapFirst(fn: (value: T) => Property): Property; - flatMapFirst>(): Property, any>; - 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, options?: { flushOnEnd?: boolean, flushOnChange?: boolean }): Property; - awaiting(otherObs: Observable): Property; +export class Stream extends Observable { + } -export interface ObservablePool extends Observable { +export class Property extends Observable { + +} + +export class Pool extends Observable { plug(obs: Observable): this; - unPlug(obs: Observable): this; + unplug(obs: Observable): this; } -export interface Event { - type: string; - value: T; -} +export type Event = + { type: 'value', value: V } | + { type: 'error', value: E } | + { type: 'end', value: void }; -export interface Emitter { - emit(value: T): void; - error(error: S): void; +export interface Emitter { + value(value: V): boolean; + event(event: Event): boolean; + error(e: E): boolean; end(): void; - emitEvent(event: { type: string, value: T | S }): void; + + // Deprecated methods + emit(value: V): boolean; + emitEvent(event: Event): boolean; } // Create a stream -export declare function never(): Stream; -export declare function later(wait: number, value: T): Stream; -export declare function interval(interval: number, value: T): Stream; -export declare function sequentially(interval: number, values: T[]): Stream; -export declare function fromPoll(interval: number, fn: () => T): Stream; -export declare function withInterval(interval: number, handler: (emitter: Emitter) => void): Stream; -export declare function fromCallback(fn: (callback: (value: T) => void) => void): Stream; -export declare function fromNodeCallback(fn: (callback: (error: S, result: T) => void) => void): Stream; -export declare function fromEvents(target: EventTarget | NodeJS.EventEmitter | { on: Function, off: Function }, eventName: string, transform?: (value: T) => S): Stream; -export declare function stream(subscribe: (emitter: Emitter) => Function | void): Stream; -export declare function fromESObservable(observable: any): 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; +export function fromESObservable(observable: any): Stream // Create a property -export declare function constant(value: T): Property; -export declare function constantError(error: T): Property; -export declare function fromPromise(promise: Promise): Property; +export function constant(value: T): Property; +export function constantError(error: T): Property; +export function fromPromise(promise: Promise): Property; // Combine observables -export declare function combine(obss: Observable[], passiveObss: Observable[], combinator?: (...values: T[]) => U): Stream; -export declare function combine(obss: Observable[], combinator: (...values: T[]) => U): Stream; -export declare function combine }>(obss: T): Stream<{ [P in keyof T]: ValueOfAnObservable }, any>; -export declare function combine], P extends keyof T>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable, Observable, Observable, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable, Observable, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable], any>; -export declare function combine]>(obss: T): Stream<[ValueOfAnObservable], any>; -export declare function combine(obss: T): Stream; -export declare function zip(obss: Observable[], passiveObss?: Observable[], combinator?: (...values: T[]) => U): Observable; -export declare function merge(obss: Observable[]): Observable; -export declare function concat(obss: Observable[]): Observable; -export declare function pool(): ObservablePool; -export declare function repeat(generator: (i: number) => Observable | boolean): Observable; +export function combine(obss: Observable[], passiveObss: Observable[], combinator?: (...values: T[]) => U): Stream; +export function combine(obss: Observable[], combinator: (...values: T[]) => U): Stream; +export function combine }>(obss: T): Stream<{ [P in keyof T]: ValueOfAnObservable }, any>; +export function combine], P extends keyof T>(obss: T): Stream<[ValueOfAnObservable], any>; +export function combine, Observable, Observable, Observable, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine, Observable, Observable, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine, Observable, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine, Observable, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine, Observable, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine, Observable, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine, Observable]>(obss: T): Stream<[ValueOfAnObservable, ValueOfAnObservable], any>; +export function combine]>(obss: T): Stream<[ValueOfAnObservable], any>; +export function combine(obss: T): Stream; +export function combine], P extends [Observable], K>(obss: T, obssP: P, combinator: (a: T[0][''], b: P[0]['']) => K): 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(): Pool; +export function repeat(generator: (i: number) => Observable | boolean): Observable; + +export var staticLand: { + Observable: { + ap(obsF: Observable<(x: A) => B, E1>, obsV: Observable): Observable; + bimap(fnE: (x: E1) => E2, fnV: (x: V1) => V2, obs: Observable): Observable; + chain(cb: (value: V) => Observable, s: Observable): Observable; + concat(obs1: Observable, obs2: Observable): Observable; + empty(): Observable; + map(cb: (value: V) => V2, s: Observable): Observable; + of(value: V): Observable; + } +} + +declare var kefir: { + Observable: typeof Observable; + Pool: typeof Pool; + Stream: typeof Stream; + Property: typeof Property; + + never: typeof never, + later: typeof later, + interval: typeof interval, + sequentially: typeof sequentially, + fromPoll: typeof fromPoll, + withInterval: typeof withInterval, + fromCallback: typeof fromCallback, + fromNodeCallback: typeof fromNodeCallback, + fromEvents: typeof fromEvents, + stream: typeof stream, + fromESObservable: typeof fromESObservable, + constant: typeof constant, + constantError: typeof constantError, + fromPromise: typeof fromPromise, + combine: typeof combine, + zip: typeof zip, + merge: typeof merge, + concat: typeof concat, + pool: typeof pool, + repeat: typeof repeat, + + staticLand: typeof staticLand +}; + +export default kefir; diff --git a/types/kefir/kefir-tests.ts b/types/kefir/kefir-tests.ts index d2dbd5835e..0fba9c5d21 100644 --- a/types/kefir/kefir-tests.ts +++ b/types/kefir/kefir-tests.ts @@ -1,6 +1,6 @@ import * as Kefir from 'kefir'; -import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'kefir'; +import { Observable, Pool, Stream, Property, Event, Emitter } from 'kefir'; //Create a stream { @@ -130,7 +130,7 @@ import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'ke 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) => { + let observable28: Stream = Kefir.sequentially(100, [0, 1, 2, 3]).withHandler((emitter: Emitter, event: Event) => { if (event.type === 'end') { emitter.emit('bye'); emitter.end(); @@ -141,6 +141,9 @@ import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'ke } } }); + type First = 'first'; + type Second = 'second'; + let observable32: Stream = Kefir.sequentially(100, ['first', 'second']).filter((value): value is First => value === 'first'); } // Combine observables @@ -177,7 +180,7 @@ import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'ke 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(); + let pool: Pool = Kefir.pool(); pool.plug(a); pool.plug(b); pool.plug(c); @@ -207,7 +210,7 @@ import { Observable, ObservablePool, Stream, Property, Event, Emitter } from 'ke { 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 observable02: Property = a.sampledBy(b) } { let foo: Stream = Kefir.sequentially(100, [1, 2, 3, 4]);