Cleanup of rxjs; removing cruft we no longer use (most was associated with the 'beta' history implementation)

This commit is contained in:
Dannon Baker
2022-08-24 16:41:23 -04:00
parent 67de2306af
commit 3237fac798
23 changed files with 23 additions and 516 deletions
+1 -1
View File
@@ -7,7 +7,7 @@
import { defer } from "rxjs";
import { shareReplay } from "rxjs/operators";
import { getGalaxyInstance } from "app";
import { waitForInit } from "utils/observable/waitForInit";
import { waitForInit } from "utils/observable";
// store subscriptions
import { syncUserToGalaxy } from "store/userStore";
@@ -1,7 +1,7 @@
// Sync Galaxy store to legacy galaxy current user
import { pluck, switchMap } from "rxjs/operators";
import { monitorBackboneModel } from "utils/observable/monitorBackboneModel";
import { monitorBackboneModel } from "utils/observable";
export function syncUserToGalaxy(galaxy$, store) {
const result$ = galaxy$.pipe(
-15
View File
@@ -1,15 +0,0 @@
/**
* Observable operator that emits one value when the source
* is emitting and another after it has stopped
*/
import { merge } from "rxjs";
import { mapTo, debounceTime, distinctUntilChanged, publish } from "rxjs/operators";
export const activity = (period = 500) => {
return publish((src) => {
const on = src.pipe(mapTo(true));
const off = src.pipe(debounceTime(period), mapTo(false));
return merge(on, off).pipe(distinctUntilChanged());
});
};
-8
View File
@@ -1,8 +0,0 @@
import { publish, buffer, debounceTime } from "rxjs/operators";
export const burst = (period = 100) => {
return publish((src) => {
const flush = src.pipe(debounceTime(period));
return src.pipe(buffer(flush));
});
};
-60
View File
@@ -1,60 +0,0 @@
import { map } from "rxjs/operators";
/**
* Bust incoming numeric source values into blocks of designated size.
*
* @param {Number|Object} cfg object or numeric chunksize
* @param {*} ceil
*/
// prettier-ignore
export const chunk = (cfg) => {
const settings = Number.isInteger(cfg) ? { chunkSize: Math.round(cfg) } : cfg;
const { chunkSize = 0, ceil = false, debug = false, label } = settings;
if (chunkSize == 0) {
throw new Error("Please provide a chunk size");
}
return map((chunkMe) => {
const rawVal = 1.0 * chunkMe / chunkSize;
const result = chunkSize * (ceil ? Math.ceil(rawVal) : Math.floor(rawVal));
if (debug) {
console.log(`chunk: ${label}`, chunkMe, result, settings);
}
return result;
})
};
/**
* Change one parameter to be an multiple of indicated block size. Used to
* regulate the URLs we send to the server so that some will be cached. Pos is
* the position in the combined inputs array, and the chunkSize is the size of
* the block to break the value into.
*
* Ex: Source Inputs [ x, y, z, 750 ],
* chunkParam(3, 200)
* Results in [ x, y, z, 600]
*
* @param {integer} pos Input array parameter number to chunk
* @param {integer} chunkSize Size of chunks
* @param {Boolean} ceil Math.ceil or Math.floor chunk value
*/
// prettier-ignore
export const chunkParam = (pos, chunkSize, ceil = false) => {
return map((inputs) => {
const chunkMe = inputs[pos];
const rawVal = 1.0 * chunkMe / chunkSize;
const chunkedVal = chunkSize * (ceil ? Math.ceil(rawVal) : Math.floor(rawVal));
const newInputs = inputs.slice();
newInputs[pos] = chunkedVal;
return newInputs;
})
}
export const chunkProp = (propName, chunkSize, ceil = false) => {
return map((obj) => {
const chunkMe = obj[propName];
const rawVal = (1.0 * chunkMe) / chunkSize;
const chunkedVal = chunkSize * (ceil ? Math.ceil(rawVal) : Math.floor(rawVal));
return { ...obj, [propName]: chunkedVal };
});
};
-21
View File
@@ -1,21 +0,0 @@
// Stateful delay gets longer with each repeated call
import { of } from "rxjs";
import { concatMap, delay } from "rxjs/operators";
import { show } from "utils/observable";
export const decay = (cfg = {}) => {
const { initialInterval = 1000, maxInterval = 60 * 1000, lambda = 0.25, debug = false } = cfg;
let counter = 0;
return concatMap((val) => {
let waitTime = Math.floor(initialInterval * Math.exp(lambda * counter++));
waitTime = Math.max(waitTime, initialInterval);
waitTime = Math.min(waitTime, maxInterval);
return of(val).pipe(
show(debug, (val) => console.log("decay", val, counter, waitTime)),
delay(waitTime)
);
});
};
@@ -1,5 +0,0 @@
import { take } from "rxjs/operators";
// take one emission, assume we're done, return as promise
// TODO: firstValueFrom will be native in rxjs soon as toPromise is being deprecated
export const firstValueFrom = (src$) => src$.pipe(take(1)).toPromise();
-16
View File
@@ -1,16 +0,0 @@
import { pipe } from "rxjs";
import { map } from "rxjs/operators";
/**
* passing objects into the worker removes class information, Pass in an array
* of constructors to be assigned, in order of the input array.
*/
// prettier-ignore
export const hydrate = (constructors = []) => pipe(
map((inputs) => {
return inputs.map((rawInput, i) => {
const C = constructors[i];
return C ? new C(rawInput) : rawInput;
});
})
)
+21 -19
View File
@@ -1,20 +1,22 @@
export { activity } from "./activity";
export { burst as debounceBurst } from "./burst";
export { chunk, chunkParam, chunkProp } from "./chunk";
export { decay } from "./decay";
export { firstValueFrom } from "./firstValueFrom";
export { hydrate } from "./hydrate";
export { monitorBackboneModel } from "./monitorBackboneModel";
export { monitorXHR } from "./monitorXHR";
export { nth } from "./nth";
export { shareButDie } from "./shareButDie";
export { show } from "./show";
export { singleton } from "./singleton";
export { throttleDistinct } from "./throttleDistinct";
export { toggle } from "./toggle";
export { waitForInit } from "./waitForInit";
export { whenAny } from "./whenAny";
export { watchVuexSelector } from "./vuex";
import { fromEvent, interval, timer } from "rxjs";
import { map, filter, take, takeUntil } from "rxjs/operators";
// Do not export rxjsDebugging in the barrel file because it has global consequences
// export { initSpy } from "./rxjsDebugging";
/**
* Creates an observable that emits a single value once it appears, as defined by the selector
* function. Emits that value and stops. Times out after designated period in the event that the
* thing never shows up.
*/
export function waitForInit(selector, cfg = {}) {
const { spamTime = 100, timeout = 10000, isValid = (val) => val !== undefined } = cfg;
return interval(spamTime).pipe(
map(() => selector()),
filter(isValid),
take(1),
takeUntil(timer(timeout))
);
}
export function monitorBackboneModel(sourceModel, evtName = "change") {
return fromEvent(sourceModel, evtName).pipe(map(([model]) => model.toJSON()));
}
@@ -1,9 +0,0 @@
import { fromEvent } from "rxjs";
import { map } from "rxjs/operators";
// prettier-ignore
export function monitorBackboneModel(sourceModel, evtName = "change") {
return fromEvent(sourceModel, evtName).pipe(
map(([model]) => model.toJSON())
);
}
-64
View File
@@ -1,64 +0,0 @@
import { defer, Subject } from "rxjs";
import { filter, finalize, share } from "rxjs/operators";
// global XHR request feed, emits every time xhr fires
const xhr$ = defer(() => {
const request$ = new Subject();
// patch
const originalOpen = XMLHttpRequest.prototype.open;
XMLHttpRequest.prototype.open = function (method, url) {
request$.next({ method, url });
return originalOpen.apply(this, arguments);
};
return request$.pipe(
// un-patch after all subscribers have left
finalize(() => {
XMLHttpRequest.prototype.open = originalOpen;
})
);
}).pipe(share());
/**
* Watch outgoing ajax calls for any of the indicated methods, filter to any of the provided route regexes
*
* @param {Object} cfg Configure methods and routes to monitor
* @return {Observable} Observable that emits every time a matching outgoing request happens
*/
export const monitorXHR = (cfg = {}) => {
return xhr$.pipe(filter(matchRoutes(cfg)));
};
// matches incoming method/url against config
export const matchRoutes =
(cfg = {}) =>
({ method, url }) => {
const { methods = ["GET", "POST", "PUT", "DELETE"], routes = [], exclude = [] } = cfg;
// match fun compares route def to current url
const matcher = matchUrlToRoute(url);
if (!methods.includes(method)) {
return false;
}
if (exclude.length && exclude.find(matcher)) {
return false;
}
if (routes.length) {
const routeMatch = routes.find(matcher);
return !!routeMatch;
}
return false;
};
const matchUrlToRoute = (url) => (rt) => {
let result = false;
if (rt instanceof RegExp) {
result = url.match(rt);
} else {
// simple matcher
result = url.includes(rt);
}
return result;
};
@@ -1,87 +0,0 @@
import { matchRoutes, monitorXHR } from "./monitorXHR";
describe("monitorXhr", () => {
describe("monitorXHR", () => {
it("should import", () => {
expect(monitorXHR).toBeInstanceOf(Function);
});
});
// request matcher, XHR object emits objects of form { method, url } which get matched to a set
// of routes of interest
describe("matchRoutes", () => {
let matcher;
beforeEach(() => {
matcher = matchRoutes({ routes: ["api/foo"] });
});
it("matcher should be a function", () => {
expect(matcher).toBeInstanceOf(Function);
});
it("should recognize an exact route", () => {
const testPacket = { method: "POST", url: "api/foo" };
expect(matcher(testPacket)).toBe(true);
});
it("should recognize a subroute", () => {
const testPacket = { method: "POST", url: "api/foo/123" };
expect(matcher(testPacket)).toBe(true);
});
it("should recognize a subroute with some query params", () => {
const testPacket = { method: "POST", url: "api/foo/123?foo=bar" };
expect(matcher(testPacket)).toBe(true);
});
it("should not match unwatched urls", () => {
const testPacket = { method: "POST", url: "api/notinthelist" };
expect(matcher(testPacket)).toBe(false);
});
it("should not match unwatched methods", () => {
const testPacket = { method: "GET", url: "api/notinthelist" };
expect(matcher(testPacket)).toBe(false);
});
});
describe("matchRoutes: general regex", () => {
it("should match foo or bar", () => {
const matcher = matchRoutes({
routes: [/api\/(foo|bar)/gi],
});
const fooTest = { method: "POST", url: "api/foo/asdfasdf" };
expect(matcher(fooTest)).toBe(true);
const barTest = { method: "POST", url: "api/bar/asdfasdf" };
expect(matcher(barTest)).toBe(true);
const blechTest = { method: "POST", url: "api/blech/asdfasdf" };
expect(matcher(blechTest)).toBe(false);
});
});
describe("matchRoutes: excludes", () => {
const matcher = matchRoutes({
routes: ["api/histories"],
exclude: ["api/histories/ABC"],
});
// should exclude urls containing the restricted ID
it("should exclude history routes with id=ABC", () => {
const testRoute = { method: "POST", url: "api/histories/ABC" };
expect(matcher(testRoute)).toBe(false);
const testRouteWithParams = { method: "POST", url: "api/histories/ABC?hoohah=123" };
expect(matcher(testRouteWithParams)).toBe(false);
});
// should pass similar routes with different ids
it("should match other history routes", () => {
const testRoute = { method: "POST", url: "api/histories/DEF" };
expect(matcher(testRoute)).toBe(true);
const testRouteWithParams = { method: "POST", url: "api/histories/DEF?foo=sadfasd" };
expect(matcher(testRouteWithParams)).toBe(true);
});
});
});
-4
View File
@@ -1,4 +0,0 @@
import { map } from "rxjs/operators";
// picks nth item of a combined array source
export const nth = (index = 0) => map((inputs) => inputs[index]);
@@ -1,31 +0,0 @@
import config from "config";
import { create } from "rxjs-spy";
import DevToolsPlugin from "rxjs-spy-devtools-plugin";
// initialize spy
let spy = undefined;
export const initSpy = () => {
if (!config.rxjsDebug || spy) {
return spy;
}
// How to use the rxjs dev panel in chrome:
// https://github.com/ardoq/rxjs-devtools/tree/master/packages/rxjs-spy-devtools-plugin
// https://github.com/cartant/rxjs-spy
console.warn("Installing rxjs spy, only do this in development");
spy = create();
const devtoolsPlugin = new DevToolsPlugin(spy, { verbose: false });
spy.plug(devtoolsPlugin);
// teardown the spy if we're hot-reloading:
if (module.hot) {
if (module.hot) {
module.hot.dispose(() => {
spy.teardown();
});
}
}
return spy;
};
@@ -1,7 +0,0 @@
/**
* Share Replay with a refcount as per the older rxjs shareReplay operator. It is often important to
* have the global availability of a shareReplay, but the ability to complete a subscription.
*/
import { shareReplay } from "rxjs/operators";
export const shareButDie = (n) => shareReplay({ bufferSize: n, refCount: true });
-13
View File
@@ -1,13 +0,0 @@
/**
* A conditional tap used for debugging so we don't have to constantly add and remove tap statements
* while working on rxjs streams
*/
import { tap } from "rxjs/operators";
export const show = (flag = false, fn) => {
return tap((val) => {
if (flag) {
fn(val);
}
});
};
-38
View File
@@ -1,38 +0,0 @@
import { of, pipe, isObservable } from "rxjs";
import { map, mergeMap } from "rxjs/operators";
// retrieves a map key for the storage from the passed input
const defaultKeyFn = (val) => val;
// stores singleton products
const defaultStorage = new Map();
/**
* Similar to map, but a unique output product value is created for each
* individual input value.
*
* @param {Function} factory Create single output from unique input
* @param {Map} storage Stores rendered product value
* @param {Function} keyFn Optional, determines storage key from passed value
*/
export const singleton = (factory, storage = defaultStorage, keyFn = defaultKeyFn) =>
pipe(
map((val) => {
const key = keyFn(val);
return storage.has(key) ? storage.get(key) : storage.set(key, factory(val)).get(key);
})
);
/**
* Generates a unique value from input source, then subscribes to that value if
* it's an observable, or renders to an observable of one value if not.
*
* @param {Function} factory Create single output from unique input
* @param {Map} storage Stores rendered product value
* @param {Function} keyFn Optional, determines storage key from passed value
*/
export const singletonMap = (factory, storage, keyFn) =>
pipe(
singleton(factory, storage, keyFn),
mergeMap((val) => (isObservable(val) ? val : of(val)))
);
@@ -1,31 +0,0 @@
import { of } from "rxjs";
import { ObserverSpy } from "@hirez_io/observer-spy";
import { singleton } from "./singleton";
describe("creating observable singleton", () => {
const source = { id: "a" };
const source$ = of(source);
// give one source observable get one product out
// each time the observable is subscribed to
test("same input should return the exact same output even from different observable streams", () => {
// build unique products with factory
const factory = jest.fn((db) => {
return { id: Math.random() };
});
const productA$ = source$.pipe(singleton(factory));
const productA2$ = source$.pipe(singleton(factory));
const spyA = new ObserverSpy();
productA$.subscribe(spyA);
const spyA2 = new ObserverSpy();
productA2$.subscribe(spyA2);
expect(spyA.getFirstValue()).toBe(spyA2.getFirstValue());
expect(factory).toHaveBeenCalledTimes(1);
});
});
@@ -1,20 +0,0 @@
import { pipe } from "rxjs";
import { groupBy, mergeMap, throttleTime } from "rxjs/operators";
/**
* Emits distinct values as they show up, but repeated values are throttled
* until the timout expires.
*
* @param {integer} timeout throttle duration
*/
export const throttleDistinct = (cfg = {}) => {
const { timeout = 1000, selector = (value) => value } = cfg;
// prettier-ignore
return pipe(
groupBy(selector),
mergeMap((grouped) => grouped.pipe(
throttleTime(timeout)
))
);
};
-14
View File
@@ -1,14 +0,0 @@
import { of, concat, EMPTY, NEVER } from "rxjs";
import { windowToggle, scan, switchAll } from "rxjs/operators";
// prettier-ignore
export const toggle = (toggleSrc$, startActive = true) => (src$) => {
const init$ = startActive ? of(1) : EMPTY;
const tgl$ = concat(init$, toggleSrc$).pipe(
scan((acc) => !acc, true)
);
return src$.pipe(
windowToggle(tgl$, (val) => (val ? src$ : NEVER)),
switchAll()
);
};
-14
View File
@@ -1,14 +0,0 @@
/**
* Function that creates on observable to watch a vuex store when provided a selector function.
*/
import { Observable } from "rxjs";
export function watchVuexSelector(store, selector, watchOptions = {}) {
const { immediate = true } = watchOptions;
return new Observable((observer) => {
const callback = (result) => observer.next(result);
return store.watch(selector, callback, { immediate });
});
}
@@ -1,29 +0,0 @@
/**
* Creates an observable that emits a single value once it appears, as defined by the selector
* function. Emits that value and stops. Times out after designated period in the event that the
* thing never shows up.
*/
import { interval, timer } from "rxjs";
import { map, filter, take, takeUntil } from "rxjs/operators";
// prettier-ignore
export function waitForInit(selector, cfg = {}) {
const {
spamTime = 100,
timeout = 10000,
isValid = (val) => val !== undefined
} = cfg;
return interval(spamTime).pipe(
map(() => selector()),
filter(isValid),
take(1),
takeUntil(timer(timeout))
);
}
export function awaitValue(selector, cfg) {
const val$ = waitForInit(selector, cfg);
return val$.toPromise();
}
-9
View File
@@ -1,9 +0,0 @@
/**
* Similar to combineLatest but adds a small debounce to avoid a duplicate event
* firing when multiple sources change at the same time.
*/
import { combineLatest } from "rxjs";
import { debounceTime } from "rxjs/operators";
export const whenAny = (...sources) => combineLatest(sources).pipe(debounceTime(0));