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