From e811d737f2e0086cbd89903e49da2bd09ab3e819 Mon Sep 17 00:00:00 2001 From: Amaury Martiny Date: Mon, 28 Jan 2019 18:46:46 +0100 Subject: [PATCH] Use memoizee instead of custom cache (#643) * First run at using memoizee * Make cached test pass * Fix replay test * Update e2e test * Fix types * Fix lint * Simplify types * Add comment --- package.json | 2 +- packages/rpc-rx/package.json | 2 + .../src/{cached.spec.js => cached.spec.ts} | 44 ++++------ .../src/{index.spec.js => index.spec.ts} | 4 +- packages/rpc-rx/src/index.ts | 69 ++++++---------- .../rpc-rx/src/{replay.js => replay.spec.ts} | 18 ++-- packages/rpc-rx/src/types.ts | 4 +- yarn.lock | 82 +++++++++++++++++-- 8 files changed, 133 insertions(+), 92 deletions(-) rename packages/rpc-rx/src/{cached.spec.js => cached.spec.ts} (55%) rename packages/rpc-rx/src/{index.spec.js => index.spec.ts} (93%) rename packages/rpc-rx/src/{replay.js => replay.spec.ts} (78%) diff --git a/package.json b/package.json index d751229041..32ef5aa73a 100644 --- a/package.json +++ b/package.json @@ -2,7 +2,7 @@ "version": "0.41.7", "private": true, "engines": { - "node": "^10.13.0", + "node": ">=10.13.0", "yarn": "^1.10.1" }, "workspaces": [ diff --git a/packages/rpc-rx/package.json b/packages/rpc-rx/package.json index 7319420345..f722667a98 100644 --- a/packages/rpc-rx/package.json +++ b/packages/rpc-rx/package.json @@ -32,7 +32,9 @@ "@babel/runtime": "^7.2.0", "@polkadot/rpc-core": "^0.41.7", "@polkadot/rpc-provider": "^0.41.7", + "@types/memoizee": "^0.4.2", "@types/rx": "^4.1.1", + "memoizee": "^0.4.14", "rxjs": "^6.3.3" } } diff --git a/packages/rpc-rx/src/cached.spec.js b/packages/rpc-rx/src/cached.spec.ts similarity index 55% rename from packages/rpc-rx/src/cached.spec.js rename to packages/rpc-rx/src/cached.spec.ts index 341a3489d1..19296b0b79 100644 --- a/packages/rpc-rx/src/cached.spec.js +++ b/packages/rpc-rx/src/cached.spec.ts @@ -2,46 +2,40 @@ // This software may be modified and distributed under the terms // of the Apache-2.0 license. See the LICENSE file for details. +import { Observable } from 'rxjs'; +import { RpcInterface$Section } from '@polkadot/rpc-core/types'; + jest.mock('@polkadot/rpc-provider/ws', () => class { isConnected = () => true; on = () => true; send = () => true; }); -const RpcRx = require('./index').default; +import RpcRx from './index'; describe('createCachedObservable', () => { - let api; - let creator; - let section; + let api: RpcRx; + let creator: (...params: Array) => Observable; + let section: RpcInterface$Section; beforeEach(() => { api = new RpcRx(); }); beforeEach(() => { - const subMethod = jest.fn((name, ...params) => { - return Promise.resolve(12345); - }); - - subMethod.unsubscribe = jest.fn(() => { - return Promise.resolve(true); - }); - - const resMethod = jest.fn((name, ...params) => { - return Promise.resolve(12345); - }); + const subMethod: any = jest.fn(() => Promise.resolve(12345)); + subMethod.unsubscribe = jest.fn(() => Promise.resolve(true)); section = { - resMethod, subMethod }; - creator = api.createCachedObservable('test', 'subMethod', section); + // @ts-ignore + creator = api.createObservable('subMethod', section); }); it('creates a single observable', () => { - creator(123).subscribe((value) => {}); + creator(123).subscribe(); expect( section.subMethod @@ -50,27 +44,17 @@ describe('createCachedObservable', () => { it('creates a single observable (multiple calls)', () => { const observable1 = creator(123); - - observable1.subscribe((value) => {}); - const observable2 = creator(123); - observable2.subscribe((value) => {}); - expect( observable2 - ).toEqual(observable1); + ).toBe(observable1); }); - it('creates multiple observers for different values', () => { + it('creates multiple observables for different values', () => { const observable1 = creator(123); - - observable1.subscribe((value) => {}); - const observable2 = creator(456); - observable2.subscribe((value) => {}); - expect( observable2 ).not.toEqual(observable1); diff --git a/packages/rpc-rx/src/index.spec.js b/packages/rpc-rx/src/index.spec.ts similarity index 93% rename from packages/rpc-rx/src/index.spec.js rename to packages/rpc-rx/src/index.spec.ts index a683e19bea..c7e3a58762 100644 --- a/packages/rpc-rx/src/index.spec.js +++ b/packages/rpc-rx/src/index.spec.ts @@ -10,10 +10,10 @@ jest.mock('@polkadot/rpc-provider/ws', () => class { send = () => true; }); -const RpcRx = require('./index').default; +import RpcRx from './index'; describe('RpcRx', () => { - let api; + let api: RpcRx; beforeEach(() => { api = new RpcRx(); diff --git a/packages/rpc-rx/src/index.ts b/packages/rpc-rx/src/index.ts index d708db526b..f82f41ea2b 100644 --- a/packages/rpc-rx/src/index.ts +++ b/packages/rpc-rx/src/index.ts @@ -4,20 +4,15 @@ import { RpcInterface$Method, RpcInterface$Section } from '@polkadot/rpc-core/types'; import { ProviderInterface } from '@polkadot/rpc-provider/types'; -import { RpcRxInterface, RpcRxInterface$Events, RpcRxInterface$Section } from './types'; +import { RpcRxInterface, RpcRxInterface$Events, RpcRxInterface$Method, RpcRxInterface$Section } from './types'; import EventEmitter from 'eventemitter3'; -import { BehaviorSubject, Observable, Subscriber, from } from 'rxjs'; -import { map, publishReplay, refCount } from 'rxjs/operators'; +import memoize, { Memoized } from 'memoizee'; +import { BehaviorSubject, Observable, from, Observer } from 'rxjs'; import Rpc from '@polkadot/rpc-core/index'; +import { map, publishReplay, refCount } from 'rxjs/operators'; import { isFunction, isUndefined } from '@polkadot/util'; -type CachedMap = { - [index: string]: { - [index: string]: Observable - } -}; - /** * @name RpcRx * @summary The RxJS API is a wrapper around the API. @@ -36,7 +31,6 @@ type CachedMap = { */ export default class RpcRx implements RpcRxInterface { private _api: Rpc; - private _cacheMap: CachedMap; private _eventemitter: EventEmitter; private _isConnected: BehaviorSubject; readonly author: RpcRxInterface$Section; @@ -51,16 +45,15 @@ export default class RpcRx implements RpcRxInterface { this._api = providerOrRpc instanceof Rpc ? providerOrRpc : new Rpc(providerOrRpc); - this._cacheMap = {}; this._eventemitter = new EventEmitter(); this._isConnected = new BehaviorSubject(this._api._provider.isConnected()); this.initEmitters(this._api._provider); - this.author = this.createInterface('author', this._api.author); - this.chain = this.createInterface('chain', this._api.chain); - this.state = this.createInterface('state', this._api.state); - this.system = this.createInterface('system', this._api.system); + this.author = this.createInterface(this._api.author); + this.chain = this.createInterface(this._api.chain); + this.state = this.createInterface(this._api.state); + this.system = this.createInterface(this._api.system); } isConnected (): BehaviorSubject { @@ -89,22 +82,30 @@ export default class RpcRx implements RpcRxInterface { }); } - private createInterface (sectionName: string, section: RpcInterface$Section): RpcRxInterface$Section { + private createInterface (section: RpcInterface$Section): RpcRxInterface$Section { return Object .keys(section) .filter((name) => !['subscribe', 'unsubscribe'].includes(name)) .reduce((observables, name) => { - observables[name] = this.createObservable(`${sectionName}_${name}`, name, section); + observables[name] = this.createObservable(name, section); return observables; }, ({} as RpcRxInterface$Section)); } - private createObservable (subName: string, name: string, section: RpcInterface$Section): (...params: Array) => Observable { + private createObservable (name: string, section: RpcInterface$Section): RpcRxInterface$Method { if (isFunction(section[name].unsubscribe)) { - return this.createCachedObservable(subName, name, section); + const memoized: Memoized = memoize( + (...params: Array) => this.createReplay(name, params, section, memoized), + { length: false } + ); + + return memoized as unknown as RpcRxInterface$Method; } + // We voluntarily don't cache the "one-shot" RPC calls. For example, + // `getStorage('123')` returns the current value, but this value can change + // over time, so we wouldn't want to cache the Observable. return (...params: Array): Observable => from( section[name] @@ -115,32 +116,16 @@ export default class RpcRx implements RpcRxInterface { ); } - private createCachedObservable (subName: string, name: string, section: RpcInterface$Section): (...params: Array) => Observable { - if (!this._cacheMap[subName]) { - this._cacheMap[subName] = {}; - } - - return (...params: Array): Observable => { - const paramStr = JSON.stringify(params); - - if (!this._cacheMap[subName][paramStr]) { - this._cacheMap[subName][paramStr] = this.createReplay(name, params, section, subName, paramStr); - } - - return this._cacheMap[subName][paramStr]; - }; - } - - private createReplay (name: string, params: Array, section: RpcInterface$Section, subName: string, paramStr: string): Observable { + private createReplay (name: string, params: Array, section: RpcInterface$Section, memoized: Memoized): Observable { return Observable - .create((observer: Subscriber): Function => { + .create((observer: Observer): Function => { const fn = section[name]; const callback = this.createReplayCallback(observer); const subscribe = fn(...params, callback).catch((error) => observer.next(error) ); - return this.createReplayUnsub(fn, subscribe, subName, paramStr); + return this.createReplayUnsub(fn, subscribe, params, memoized); }) .pipe( map((value) => { @@ -155,12 +140,12 @@ export default class RpcRx implements RpcRxInterface { ); } - private createReplayCallback (observer: Subscriber) { + private createReplayCallback (observer: Observer) { let cachedResult: any; return (result: any) => { if (isUndefined(cachedResult) || !Array.isArray(cachedResult) || !Array.isArray(result) - || result.length !== cachedResult.length) { + || result.length !== cachedResult.length) { cachedResult = result; } else { cachedResult = cachedResult.map((cachedValue, index) => @@ -174,14 +159,14 @@ export default class RpcRx implements RpcRxInterface { }; } - private createReplayUnsub (fn: RpcInterface$Method, subscribe: Promise, subName: string, paramStr: string): () => void { + private createReplayUnsub (fn: RpcInterface$Method, subscribe: Promise, params: Array, memoized: Memoized): () => void { return (): void => { subscribe .then((subscriptionId: number) => fn.unsubscribe(subscriptionId) ) .then(() => { - delete this._cacheMap[subName][paramStr]; + memoized.delete(...params); }) .catch((error) => { console.error('Unsubscribe failed', error); diff --git a/packages/rpc-rx/src/replay.js b/packages/rpc-rx/src/replay.spec.ts similarity index 78% rename from packages/rpc-rx/src/replay.js rename to packages/rpc-rx/src/replay.spec.ts index 193cd48c04..717b441e6d 100644 --- a/packages/rpc-rx/src/replay.js +++ b/packages/rpc-rx/src/replay.spec.ts @@ -2,27 +2,30 @@ // This software may be modified and distributed under the terms // of the Apache-2.0 license. See the LICENSE file for details. +import { Observable } from 'rxjs'; +import { RpcInterface$Section } from '@polkadot/rpc-core/types'; + jest.mock('@polkadot/rpc-provider/ws', () => class { isConnected = () => true; on = () => true; send = () => true; }); -const RpcRx = require('./index').default; +import RpcRx from './index'; describe('replay', () => { const params = [123, false]; - let api; - let update; - let section; - let observable; + let api: RpcRx; + let section: RpcInterface$Section; + let observable: Observable; + let update: any; beforeEach(() => { api = new RpcRx(); }); beforeEach(() => { - const subMethod = jest.fn((name, ...params) => { + const subMethod: any = jest.fn((name, ...params) => { update = params.pop(); return Promise.resolve(12345); @@ -36,7 +39,8 @@ describe('replay', () => { subMethod }; - observable = api.createReplay('subMethod', params, section); + // @ts-ignore + observable = api.createObservable('subMethod', section)(...params); }); it('subscribes via the api section', (done) => { diff --git a/packages/rpc-rx/src/types.ts b/packages/rpc-rx/src/types.ts index 60328d4bdf..b0e93399f0 100644 --- a/packages/rpc-rx/src/types.ts +++ b/packages/rpc-rx/src/types.ts @@ -2,10 +2,10 @@ // This software may be modified and distributed under the terms // of the Apache-2.0 license. See the LICENSE file for details. -import { ReplaySubject, Observable } from 'rxjs'; +import { Observable } from 'rxjs'; import { ProviderInterface$Emitted } from '@polkadot/rpc-provider/types'; -export type RpcRxInterface$Method = (...params: Array) => Observable | ReplaySubject; +export type RpcRxInterface$Method = (...params: Array) => Observable; export type RpcRxInterface$Section = { [index: string]: RpcRxInterface$Method diff --git a/yarn.lock b/yarn.lock index 0b878f9f80..286ce5a6e6 100644 --- a/yarn.lock +++ b/yarn.lock @@ -2003,12 +2003,37 @@ babel-code-frame@^6.22.0, babel-code-frame@^6.26.0: esutils "^2.0.2" js-tokens "^3.0.2" -babel-core@^6.0.0, babel-core@^7.0.0-bridge.0: +babel-core@^6.0.0, babel-core@^6.26.0: + version "6.26.3" + resolved "https://registry.yarnpkg.com/babel-core/-/babel-core-6.26.3.tgz#b2e2f09e342d0f0c88e2f02e067794125e75c207" + integrity sha512-6jyFLuDmeidKmUEb3NM+/yawG0M2bDZ9Z1qbZP59cyHLz8kYGKYwpJP0UwUKKUiTRNvxfLesJnTedqczP7cTDA== + dependencies: + babel-code-frame "^6.26.0" + babel-generator "^6.26.0" + babel-helpers "^6.24.1" + babel-messages "^6.23.0" + babel-register "^6.26.0" + babel-runtime "^6.26.0" + babel-template "^6.26.0" + babel-traverse "^6.26.0" + babel-types "^6.26.0" + babylon "^6.18.0" + convert-source-map "^1.5.1" + debug "^2.6.9" + json5 "^0.5.1" + lodash "^4.17.4" + minimatch "^3.0.4" + path-is-absolute "^1.0.1" + private "^0.1.8" + slash "^1.0.0" + source-map "^0.5.7" + +babel-core@^7.0.0-bridge.0: version "7.0.0-bridge.0" resolved "https://registry.yarnpkg.com/babel-core/-/babel-core-7.0.0-bridge.0.tgz#95a492ddd90f9b4e9a4a1da14eb335b87b634ece" integrity sha512-poPX9mZH/5CSanm50Q+1toVci6pv5KSRv/5TWCwtzQS5XEwn40BcCrgIeMFWP9CKKIniKXNxoIOnOq4VVlGXhg== -babel-generator@^6.18.0: +babel-generator@^6.18.0, babel-generator@^6.26.0: version "6.26.1" resolved "https://registry.yarnpkg.com/babel-generator/-/babel-generator-6.26.1.tgz#1844408d3b8f0d35a404ea7ac180f087a601bd90" integrity sha512-HyfwY6ApZj7BYTcJURpM5tznulaBvyio7/0d4zFOeMPUmfxkCjHocCuoLa2SAGzBI8AREcH3eP3758F672DppA== @@ -2022,6 +2047,14 @@ babel-generator@^6.18.0: source-map "^0.5.7" trim-right "^1.0.1" +babel-helpers@^6.24.1: + version "6.24.1" + resolved "https://registry.yarnpkg.com/babel-helpers/-/babel-helpers-6.24.1.tgz#3471de9caec388e5c850e597e58a26ddf37602b2" + integrity sha1-NHHenK7DiOXIUOWX5Yom3fN2ArI= + dependencies: + babel-runtime "^6.22.0" + babel-template "^6.24.1" + babel-jest@^23.6.0: version "23.6.0" resolved "https://registry.yarnpkg.com/babel-jest/-/babel-jest-23.6.0.tgz#a644232366557a2240a0c083da6b25786185a2f1" @@ -2076,6 +2109,19 @@ babel-preset-jest@^23.2.0: babel-plugin-jest-hoist "^23.2.0" babel-plugin-syntax-object-rest-spread "^6.13.0" +babel-register@^6.26.0: + version "6.26.0" + resolved "https://registry.yarnpkg.com/babel-register/-/babel-register-6.26.0.tgz#6ed021173e2fcb486d7acb45c6009a856f647071" + integrity sha1-btAhFz4vy0htestFxgCahW9kcHE= + dependencies: + babel-core "^6.26.0" + babel-runtime "^6.26.0" + core-js "^2.5.0" + home-or-tmp "^2.0.0" + lodash "^4.17.4" + mkdirp "^0.5.1" + source-map-support "^0.4.15" + babel-runtime@^6.22.0, babel-runtime@^6.26.0, babel-runtime@^6.9.2: version "6.26.0" resolved "https://registry.yarnpkg.com/babel-runtime/-/babel-runtime-6.26.0.tgz#965c7058668e82b55d7bfe04ff2337bc8b5647fe" @@ -2084,7 +2130,7 @@ babel-runtime@^6.22.0, babel-runtime@^6.26.0, babel-runtime@^6.9.2: core-js "^2.4.0" regenerator-runtime "^0.11.0" -babel-template@^6.16.0: +babel-template@^6.16.0, babel-template@^6.24.1, babel-template@^6.26.0: version "6.26.0" resolved "https://registry.yarnpkg.com/babel-template/-/babel-template-6.26.0.tgz#de03e2d16396b069f46dd9fff8521fb1a0e35e02" integrity sha1-3gPi0WOWsGn0bdn/+FIfsaDjXgI= @@ -2949,7 +2995,7 @@ conventional-recommended-bump@^4.0.4: meow "^4.0.0" q "^1.5.1" -convert-source-map@^1.1.0, convert-source-map@^1.4.0: +convert-source-map@^1.1.0, convert-source-map@^1.4.0, convert-source-map@^1.5.1: version "1.6.0" resolved "https://registry.yarnpkg.com/convert-source-map/-/convert-source-map-1.6.0.tgz#51b537a8c43e0f04dec1993bffcdd504e758ac20" integrity sha512-eFu7XigvxdZ1ETfbgPBohgyQ/Z++C0eEhTor0qRwBw9unw+L0/6V8wkSuGgzdThkiS5lSpdptOQPD8Ak40a+7A== @@ -2978,6 +3024,11 @@ core-js@^2.4.0, core-js@^2.5.7: resolved "https://registry.yarnpkg.com/core-js/-/core-js-2.6.1.tgz#87416ae817de957a3f249b3b5ca475d4aaed6042" integrity sha512-L72mmmEayPJBejKIWe2pYtGis5r0tQ5NaJekdhyXgeMQTpJoBsH0NL4ElY2LfSoV15xeQWKQ+XTTOZdyero5Xg== +core-js@^2.5.0: + version "2.6.3" + resolved "https://registry.yarnpkg.com/core-js/-/core-js-2.6.3.tgz#4b70938bdffdaf64931e66e2db158f0892289c49" + integrity sha512-l00tmFFZOBHtYhN4Cz7k32VM7vTn3rE2ANjQDxdEN6zmXZ/xq1jQuutnmHvMG1ZJ7xd72+TA5YpUK8wz3rWsfQ== + core-js@^2.6.2: version "2.6.2" resolved "https://registry.yarnpkg.com/core-js/-/core-js-2.6.2.tgz#267988d7268323b349e20b4588211655f0e83944" @@ -3203,7 +3254,7 @@ debug@3.1.0: dependencies: ms "2.0.0" -debug@^2.1.2, debug@^2.2.0, debug@^2.3.3, debug@^2.6.8: +debug@^2.1.2, debug@^2.2.0, debug@^2.3.3, debug@^2.6.8, debug@^2.6.9: version "2.6.9" resolved "https://registry.yarnpkg.com/debug/-/debug-2.6.9.tgz#5d128515df134ff327e90a4c93f4e077a536341f" integrity sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA== @@ -4559,6 +4610,14 @@ hoek@2.x.x: resolved "https://registry.yarnpkg.com/hoek/-/hoek-2.16.3.tgz#20bb7403d3cea398e91dc4710a8ff1b8274a25ed" integrity sha1-ILt0A9POo5jpHcRxCo/xuCdKJe0= +home-or-tmp@^2.0.0: + version "2.0.0" + resolved "https://registry.yarnpkg.com/home-or-tmp/-/home-or-tmp-2.0.0.tgz#e36c3f2d2cae7d746a857e38d18d5f32a7882db8" + integrity sha1-42w/LSyufXRqhX440Y1fMqeILbg= + dependencies: + os-homedir "^1.0.0" + os-tmpdir "^1.0.1" + home-or-tmp@^3.0.0: version "3.0.0" resolved "https://registry.yarnpkg.com/home-or-tmp/-/home-or-tmp-3.0.0.tgz#57a8fe24cf33cdd524860a15821ddc25c86671fb" @@ -7251,7 +7310,7 @@ os-locale@^3.0.0: lcid "^2.0.0" mem "^4.0.0" -os-tmpdir@^1.0.0, os-tmpdir@~1.0.1, os-tmpdir@~1.0.2: +os-tmpdir@^1.0.0, os-tmpdir@^1.0.1, os-tmpdir@~1.0.1, os-tmpdir@~1.0.2: version "1.0.2" resolved "https://registry.yarnpkg.com/os-tmpdir/-/os-tmpdir-1.0.2.tgz#bbe67406c79aa85c5cfec766fe5734555dfa1274" integrity sha1-u+Z0BseaqFxc/sdm/lc0VV36EnQ= @@ -7496,7 +7555,7 @@ path-exists@^3.0.0: resolved "https://registry.yarnpkg.com/path-exists/-/path-exists-3.0.0.tgz#ce0ebeaa5f78cb18925ea7d810d7b59b010fd515" integrity sha1-zg6+ql94yxiSXqfYENe1mwEP1RU= -path-is-absolute@^1.0.0: +path-is-absolute@^1.0.0, path-is-absolute@^1.0.1: version "1.0.1" resolved "https://registry.yarnpkg.com/path-is-absolute/-/path-is-absolute-1.0.1.tgz#174b9268735534ffbc7ace6bf53a5a9e1b5c5f5f" integrity sha1-F0uSaHNVNP+8es5r9TpanhtcX18= @@ -7634,7 +7693,7 @@ pretty-format@^23.6.0: ansi-regex "^3.0.0" ansi-styles "^3.2.0" -private@^0.1.6: +private@^0.1.6, private@^0.1.8: version "0.1.8" resolved "https://registry.yarnpkg.com/private/-/private-0.1.8.tgz#2381edb3689f7a53d653190060fcf822d2f368ff" integrity sha512-VvivMrbvd2nKkiG38qjULzlc+4Vx4wm/whI9pQD35YrARNnhxeiRktSOhSukRLFNlzg6Br/cJPet5J/u19r/mg== @@ -8630,6 +8689,13 @@ source-map-resolve@^0.5.0: source-map-url "^0.4.0" urix "^0.1.0" +source-map-support@^0.4.15: + version "0.4.18" + resolved "https://registry.yarnpkg.com/source-map-support/-/source-map-support-0.4.18.tgz#0286a6de8be42641338594e97ccea75f0a2c585f" + integrity sha512-try0/JqxPLF9nOjvSta7tVondkP5dwgyLDjVoyMDlmjugT2lRZ1OfsrYTkCd2hkDnJTKRbO/Rl3orm8vlsUzbA== + dependencies: + source-map "^0.5.6" + source-map-support@^0.5.6, source-map-support@^0.5.9: version "0.5.9" resolved "https://registry.yarnpkg.com/source-map-support/-/source-map-support-0.5.9.tgz#41bc953b2534267ea2d605bccfa7bfa3111ced5f"