From 3237fac798a8577918ed669cea85c8246e2779bd Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Wed, 24 Aug 2022 16:40:30 -0400 Subject: [PATCH] Cleanup of rxjs; removing cruft we no longer use (most was associated with the 'beta' history implementation) --- client/src/store/syncVuextoGalaxy.js | 2 +- .../src/store/userStore/syncUserToGalaxy.js | 2 +- client/src/utils/observable/activity.js | 15 ---- client/src/utils/observable/burst.js | 8 -- client/src/utils/observable/chunk.js | 60 ------------- client/src/utils/observable/decay.js | 21 ----- client/src/utils/observable/firstValueFrom.js | 5 -- client/src/utils/observable/hydrate.js | 16 ---- client/src/utils/observable/index.js | 40 +++++---- .../utils/observable/monitorBackboneModel.js | 9 -- client/src/utils/observable/monitorXHR.js | 64 -------------- .../src/utils/observable/monitorXHR.test.js | 87 ------------------- client/src/utils/observable/nth.js | 4 - client/src/utils/observable/rxjsDebugging.js | 31 ------- client/src/utils/observable/shareButDie.js | 7 -- client/src/utils/observable/show.js | 13 --- client/src/utils/observable/singleton.js | 38 -------- client/src/utils/observable/singleton.test.js | 31 ------- .../src/utils/observable/throttleDistinct.js | 20 ----- client/src/utils/observable/toggle.js | 14 --- client/src/utils/observable/vuex.js | 14 --- client/src/utils/observable/waitForInit.js | 29 ------- client/src/utils/observable/whenAny.js | 9 -- 23 files changed, 23 insertions(+), 516 deletions(-) delete mode 100644 client/src/utils/observable/activity.js delete mode 100644 client/src/utils/observable/burst.js delete mode 100644 client/src/utils/observable/chunk.js delete mode 100644 client/src/utils/observable/decay.js delete mode 100644 client/src/utils/observable/firstValueFrom.js delete mode 100644 client/src/utils/observable/hydrate.js delete mode 100644 client/src/utils/observable/monitorBackboneModel.js delete mode 100644 client/src/utils/observable/monitorXHR.js delete mode 100644 client/src/utils/observable/monitorXHR.test.js delete mode 100644 client/src/utils/observable/nth.js delete mode 100644 client/src/utils/observable/rxjsDebugging.js delete mode 100644 client/src/utils/observable/shareButDie.js delete mode 100644 client/src/utils/observable/show.js delete mode 100644 client/src/utils/observable/singleton.js delete mode 100644 client/src/utils/observable/singleton.test.js delete mode 100644 client/src/utils/observable/throttleDistinct.js delete mode 100644 client/src/utils/observable/toggle.js delete mode 100644 client/src/utils/observable/vuex.js delete mode 100644 client/src/utils/observable/waitForInit.js delete mode 100644 client/src/utils/observable/whenAny.js diff --git a/client/src/store/syncVuextoGalaxy.js b/client/src/store/syncVuextoGalaxy.js index 5ea215d5b56..b173ecd8e51 100644 --- a/client/src/store/syncVuextoGalaxy.js +++ b/client/src/store/syncVuextoGalaxy.js @@ -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"; diff --git a/client/src/store/userStore/syncUserToGalaxy.js b/client/src/store/userStore/syncUserToGalaxy.js index d13ff7a211b..3ab06eb791c 100644 --- a/client/src/store/userStore/syncUserToGalaxy.js +++ b/client/src/store/userStore/syncUserToGalaxy.js @@ -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( diff --git a/client/src/utils/observable/activity.js b/client/src/utils/observable/activity.js deleted file mode 100644 index 5e6d0fd0276..00000000000 --- a/client/src/utils/observable/activity.js +++ /dev/null @@ -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()); - }); -}; diff --git a/client/src/utils/observable/burst.js b/client/src/utils/observable/burst.js deleted file mode 100644 index 48ee8be86eb..00000000000 --- a/client/src/utils/observable/burst.js +++ /dev/null @@ -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)); - }); -}; diff --git a/client/src/utils/observable/chunk.js b/client/src/utils/observable/chunk.js deleted file mode 100644 index 850d6358705..00000000000 --- a/client/src/utils/observable/chunk.js +++ /dev/null @@ -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 }; - }); -}; diff --git a/client/src/utils/observable/decay.js b/client/src/utils/observable/decay.js deleted file mode 100644 index 06c5cf37b3c..00000000000 --- a/client/src/utils/observable/decay.js +++ /dev/null @@ -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) - ); - }); -}; diff --git a/client/src/utils/observable/firstValueFrom.js b/client/src/utils/observable/firstValueFrom.js deleted file mode 100644 index b34391ebdec..00000000000 --- a/client/src/utils/observable/firstValueFrom.js +++ /dev/null @@ -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(); diff --git a/client/src/utils/observable/hydrate.js b/client/src/utils/observable/hydrate.js deleted file mode 100644 index 01212726423..00000000000 --- a/client/src/utils/observable/hydrate.js +++ /dev/null @@ -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; - }); - }) -) diff --git a/client/src/utils/observable/index.js b/client/src/utils/observable/index.js index 430e19c50ff..9f4c80dee49 100644 --- a/client/src/utils/observable/index.js +++ b/client/src/utils/observable/index.js @@ -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())); +} diff --git a/client/src/utils/observable/monitorBackboneModel.js b/client/src/utils/observable/monitorBackboneModel.js deleted file mode 100644 index 8f13378efde..00000000000 --- a/client/src/utils/observable/monitorBackboneModel.js +++ /dev/null @@ -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()) - ); -} diff --git a/client/src/utils/observable/monitorXHR.js b/client/src/utils/observable/monitorXHR.js deleted file mode 100644 index 54a71470ed8..00000000000 --- a/client/src/utils/observable/monitorXHR.js +++ /dev/null @@ -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; -}; diff --git a/client/src/utils/observable/monitorXHR.test.js b/client/src/utils/observable/monitorXHR.test.js deleted file mode 100644 index fd8274e1240..00000000000 --- a/client/src/utils/observable/monitorXHR.test.js +++ /dev/null @@ -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); - }); - }); -}); diff --git a/client/src/utils/observable/nth.js b/client/src/utils/observable/nth.js deleted file mode 100644 index 24496524015..00000000000 --- a/client/src/utils/observable/nth.js +++ /dev/null @@ -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]); diff --git a/client/src/utils/observable/rxjsDebugging.js b/client/src/utils/observable/rxjsDebugging.js deleted file mode 100644 index 74f6fcd70ae..00000000000 --- a/client/src/utils/observable/rxjsDebugging.js +++ /dev/null @@ -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; -}; diff --git a/client/src/utils/observable/shareButDie.js b/client/src/utils/observable/shareButDie.js deleted file mode 100644 index 91aaf882154..00000000000 --- a/client/src/utils/observable/shareButDie.js +++ /dev/null @@ -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 }); diff --git a/client/src/utils/observable/show.js b/client/src/utils/observable/show.js deleted file mode 100644 index 7cda05cc632..00000000000 --- a/client/src/utils/observable/show.js +++ /dev/null @@ -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); - } - }); -}; diff --git a/client/src/utils/observable/singleton.js b/client/src/utils/observable/singleton.js deleted file mode 100644 index 4af654acbf6..00000000000 --- a/client/src/utils/observable/singleton.js +++ /dev/null @@ -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))) - ); diff --git a/client/src/utils/observable/singleton.test.js b/client/src/utils/observable/singleton.test.js deleted file mode 100644 index c143949acb4..00000000000 --- a/client/src/utils/observable/singleton.test.js +++ /dev/null @@ -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); - }); -}); diff --git a/client/src/utils/observable/throttleDistinct.js b/client/src/utils/observable/throttleDistinct.js deleted file mode 100644 index 107cf78ade2..00000000000 --- a/client/src/utils/observable/throttleDistinct.js +++ /dev/null @@ -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) - )) - ); -}; diff --git a/client/src/utils/observable/toggle.js b/client/src/utils/observable/toggle.js deleted file mode 100644 index 3b467d9b168..00000000000 --- a/client/src/utils/observable/toggle.js +++ /dev/null @@ -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() - ); -}; diff --git a/client/src/utils/observable/vuex.js b/client/src/utils/observable/vuex.js deleted file mode 100644 index 9aa115fdc71..00000000000 --- a/client/src/utils/observable/vuex.js +++ /dev/null @@ -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 }); - }); -} diff --git a/client/src/utils/observable/waitForInit.js b/client/src/utils/observable/waitForInit.js deleted file mode 100644 index f104a12856c..00000000000 --- a/client/src/utils/observable/waitForInit.js +++ /dev/null @@ -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(); -} diff --git a/client/src/utils/observable/whenAny.js b/client/src/utils/observable/whenAny.js deleted file mode 100644 index c79e13f8a3d..00000000000 --- a/client/src/utils/observable/whenAny.js +++ /dev/null @@ -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));