diff --git a/client/src/components/History/providers/HistoryContentProvider/loadContents.js b/client/src/components/History/providers/HistoryContentProvider/loadContents.js index a18488e4028..8a57785f2e8 100644 --- a/client/src/components/History/providers/HistoryContentProvider/loadContents.js +++ b/client/src/components/History/providers/HistoryContentProvider/loadContents.js @@ -1,7 +1,7 @@ -import { pipe, defer, of } from "rxjs"; -import { repeat, switchMap, catchError } from "rxjs/operators"; -import { tag } from "rxjs-spy/operators/tag"; -import { decay } from "../../../../utils/observable/decay"; +import { defer, of } from "rxjs"; +import { repeat, switchMap, switchMapTo, startWith, debounceTime } from "rxjs/operators"; +import { decay } from "utils/observable/decay"; +import { monitorXHR } from "utils/observable/monitorXHR"; import { loadHistoryContents } from "../../caching"; // prettier-ignore @@ -13,21 +13,32 @@ export const loadContents = (cfg = {}) => { disablePoll = false, } = cfg; - return pipe( - switchMap(([{id}, params, hid]) => { - const singleLoad$ = defer(() => of([id, params, hid]).pipe( - tag('ajaxLoadInputs'), - loadHistoryContents({ windowSize }), - )); - const poll$ = singleLoad$.pipe( - decay({ initialInterval, maxInterval }), - repeat(), - ); - return disablePoll ? singleLoad$ : poll$; - }), - catchError(err => { - console.warn("Error in loadContents", err); - throw err; - }) - ); + return switchMap(([{id}, params, hid]) => { + + // a single history update + const singleLoad$ = defer(() => of([id, params, hid]).pipe( + loadHistoryContents({ windowSize }), + )); + + // start repeating, delay gets longer over time until unsubscribed + const freshPoll$ = singleLoad$.pipe( + decay({ initialInterval, maxInterval }), + repeat(), + ); + + // history, tools routes all refresh + // exclude our own polling url though + const routes = [/api\/(history|tools|histories)/]; + const methods = ["POST", "PUT", "DELETE"] + const resetPoll$ = monitorXHR({ methods, routes }); + + // resets re-subscribe to freshPoll$ starting the decay over again + const poll$ = resetPoll$.pipe( + startWith(true), + debounceTime(100), + switchMapTo(freshPoll$) + ); + + return disablePoll ? singleLoad$ : poll$; + }) } diff --git a/client/src/utils/observable/decay.js b/client/src/utils/observable/decay.js index 5d297b00d78..1c3b702da7c 100644 --- a/client/src/utils/observable/decay.js +++ b/client/src/utils/observable/decay.js @@ -1,5 +1,5 @@ // Stateful delay gets longer with each repeated call -import { pipe, of } from "rxjs"; +import { of } from "rxjs"; import { mergeMap, delay } from "rxjs/operators"; // prettier-ignore @@ -16,13 +16,11 @@ export const decay = (cfg = {}) => { let counter = 0; - return pipe( - mergeMap(val => { - let waitTime = Math.floor(initialInterval * Math.exp(lambda * counter++)); - waitTime = Math.max(waitTime, initialInterval); - waitTime = Math.min(waitTime, maxInterval); - // console.log("decay time", waitTime); - return of(val).pipe(delay(waitTime)) - }) - ); + return mergeMap(val => { + let waitTime = Math.floor(initialInterval * Math.exp(lambda * counter++)); + waitTime = Math.max(waitTime, initialInterval); + waitTime = Math.min(waitTime, maxInterval); + // console.log("decay time", waitTime); + return of(val).pipe(delay(waitTime)) + }); } diff --git a/client/src/utils/observable/monitorXHR.js b/client/src/utils/observable/monitorXHR.js new file mode 100644 index 00000000000..74c2128ed6a --- /dev/null +++ b/client/src/utils/observable/monitorXHR.js @@ -0,0 +1,65 @@ +import { defer, Subject } from "rxjs"; +import { filter, finalize, share } from "rxjs/operators"; + +// global XHR request feed, emits every time xhr fires +// prettier-ignore +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 outgong 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 new file mode 100644 index 00000000000..fd8274e1240 --- /dev/null +++ b/client/src/utils/observable/monitorXHR.test.js @@ -0,0 +1,87 @@ +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); + }); + }); +});