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
This commit is contained in:
Amaury Martiny
2019-01-28 18:46:46 +01:00
committed by GitHub
parent eccaa75192
commit e811d737f2
8 changed files with 133 additions and 92 deletions
+1 -1
View File
@@ -2,7 +2,7 @@
"version": "0.41.7",
"private": true,
"engines": {
"node": "^10.13.0",
"node": ">=10.13.0",
"yarn": "^1.10.1"
},
"workspaces": [
+2
View File
@@ -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"
}
}
@@ -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<any>) => Observable<any>;
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);
@@ -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();
+27 -42
View File
@@ -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<any>
}
};
/**
* @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<boolean>;
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<boolean> {
@@ -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<any>) => Observable<any> {
private createObservable (name: string, section: RpcInterface$Section): RpcRxInterface$Method {
if (isFunction(section[name].unsubscribe)) {
return this.createCachedObservable(subName, name, section);
const memoized: Memoized<RpcRxInterface$Method> = memoize(
(...params: Array<any>) => 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<any>): Observable<any> =>
from(
section[name]
@@ -115,32 +116,16 @@ export default class RpcRx implements RpcRxInterface {
);
}
private createCachedObservable (subName: string, name: string, section: RpcInterface$Section): (...params: Array<any>) => Observable<any> {
if (!this._cacheMap[subName]) {
this._cacheMap[subName] = {};
}
return (...params: Array<any>): Observable<any> => {
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<any>, section: RpcInterface$Section, subName: string, paramStr: string): Observable<any> {
private createReplay (name: string, params: Array<any>, section: RpcInterface$Section, memoized: Memoized<RpcRxInterface$Method>): Observable<any> {
return Observable
.create((observer: Subscriber<any>): Function => {
.create((observer: Observer<any>): 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<any>) {
private createReplayCallback (observer: Observer<any>) {
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<number>, subName: string, paramStr: string): () => void {
private createReplayUnsub (fn: RpcInterface$Method, subscribe: Promise<number>, params: Array<any>, memoized: Memoized<RpcRxInterface$Method>): () => 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);
@@ -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<any>;
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) => {
+2 -2
View File
@@ -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<any>) => Observable<any> | ReplaySubject<any>;
export type RpcRxInterface$Method = (...params: Array<any>) => Observable<any>;
export type RpcRxInterface$Section = {
[index: string]: RpcRxInterface$Method
+74 -8
View File
@@ -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"