import { checkFunction } from "../_check" import { Source } from "../_core" import { makeObservable } from "../_obs" import { Transaction } from "../_tx" import { curry2 } from "../_util" import { dispatcherOf, Observable } from "../Observable" import { Identity } from "./_base" import { map } from "./map" export interface DoActionOp { (f: (val: T) => void, observable: Observable): Observable (f: (val: T) => void): (observable: Observable) => Observable } export interface DoErrorOp { (f: (err: Error) => void, observable: Observable): Observable (f: (err: Error) => void): (observable: Observable) => Observable } export interface DoEndOp { (f: () => void, observable: Observable): Observable (f: () => void): (observable: Observable) => Observable } export interface DoLogOp { (label: string | undefined, observable: Observable): Observable (label: string | undefined): (observable: Observable) => Observable } export const doAction: DoActionOp = curry2(_doAction) export const doError: DoErrorOp = curry2(_doError) export const doEnd: DoEndOp = curry2(_doEnd) export const doLog: DoLogOp = curry2(_doLog) function _doAction(f: (val: T) => void, observable: Observable): Observable { checkFunction(f) const eff = (val: T): T => { f(val) return val } return map(eff, observable) } function _doError(f: (err: Error) => void, observable: Observable): Observable { checkFunction(f) return makeObservable(observable, new DoError(dispatcherOf(observable), f)) } function _doEnd(f: () => void, observable: Observable): Observable { checkFunction(f) return makeObservable(observable, new DoEnd(dispatcherOf(observable), f)) } function _doLog(label: string | undefined, observable: Observable): Observable { return makeObservable(observable, new DoLog(dispatcherOf(observable), label)) } class DoError extends Identity { constructor(src: Source, private f: (err: Error) => void) { super(src) } public error(tx: Transaction, err: Error): void { const { f } = this f(err) this.sink.error(tx, err) } } class DoEnd extends Identity { constructor(src: Source, private f: () => void) { super(src) } public end(tx: Transaction): void { const { f } = this f() this.sink.end(tx) } } class DoLog extends Identity { constructor(src: Source, private label: string | undefined) { super(src) } public next(tx: Transaction, val: T): void { this.log(val) this.sink.next(tx, val) } public error(tx: Transaction, err: Error): void { this.log("", err) this.sink.error(tx, err) } public end(tx: Transaction): void { this.log("") this.sink.end(tx) } private log(...msgs: any[]): void { const args: any[] = this.label === undefined ? msgs : [this.label, ...msgs] // tslint:disable-next-line:no-console console.log(...args) } }